< Summary

Information
Class: Ice.ConnectionI.OutgoingMessage
Assembly: Ice
File(s): /_/csharp/src/Ice/ConnectionI.cs
Tag: 125_37167941578
Line coverage
100%
Covered lines: 22
Uncovered lines: 0
Coverable lines: 22
Total lines: 2970
Line coverage: 100%
Branch coverage
100%
Covered branches: 8
Total branches: 8
Branch coverage: 100%
Method coverage
100%
Covered methods: 5
Fully covered methods: 5
Total methods: 5
Method coverage: 100%
Full method coverage: 100%

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
.ctor(...)100%11100%
.ctor(...)100%11100%
canceled()100%11100%
sent()100%44100%
completed(...)100%44100%

File(s)

/_/csharp/src/Ice/ConnectionI.cs

#LineLine coverage
 1// Copyright (c) ZeroC, Inc.
 2
 3using Ice.Instrumentation;
 4using Ice.Internal;
 5using System.Diagnostics;
 6using System.Text;
 7
 8namespace Ice;
 9
 10#pragma warning disable CA1001 // _inactivityTimer is disposed by cancelInactivityTimer.
 11public sealed class ConnectionI : Internal.EventHandler, CancellationHandler, Connection
 12#pragma warning restore CA1001
 13{
 14    internal interface StartCallback
 15    {
 16        void connectionStartCompleted(ConnectionI connection);
 17
 18        void connectionStartFailed(ConnectionI connection, LocalException ex);
 19    }
 20
 21    internal void start(StartCallback callback)
 22    {
 23        try
 24        {
 25            lock (_mutex)
 26            {
 27                //
 28                // The connection might already be closed if the communicator was destroyed.
 29                //
 30                if (_state >= StateClosed)
 31                {
 32                    Debug.Assert(_exception is not null);
 33                    throw _exception;
 34                }
 35
 36                if (!initialize(SocketOperation.None) || !validate(SocketOperation.None))
 37                {
 38                    if (_connectTimeout > TimeSpan.Zero)
 39                    {
 40#pragma warning disable CA2000 // connectTimer is disposed by connectTimedOut.
 41                        var connectTimer = new System.Threading.Timer(
 42                            timerObj => connectTimedOut((System.Threading.Timer)timerObj));
 43                        // schedule timer to run once; connectTimedOut disposes the timer too.
 44                        connectTimer.Change(_connectTimeout, Timeout.InfiniteTimeSpan);
 45#pragma warning restore CA2000
 46                    }
 47
 48                    _startCallback = callback;
 49                    return;
 50                }
 51
 52                // The connection starts in the holding state. It will be activated by the connection factory.
 53                setState(StateHolding);
 54            }
 55        }
 56        catch (LocalException ex)
 57        {
 58            exception(ex);
 59            callback.connectionStartFailed(this, _exception);
 60            return;
 61        }
 62
 63        callback.connectionStartCompleted(this);
 64    }
 65
 66    internal void startAndWait()
 67    {
 68        try
 69        {
 70            lock (_mutex)
 71            {
 72                //
 73                // The connection might already be closed if the communicator was destroyed.
 74                //
 75                if (_state >= StateClosed)
 76                {
 77                    Debug.Assert(_exception is not null);
 78                    throw _exception;
 79                }
 80
 81                if (!initialize(SocketOperation.None) || !validate(SocketOperation.None))
 82                {
 83                    //
 84                    // Wait for the connection to be validated.
 85                    //
 86                    while (_state <= StateNotValidated)
 87                    {
 88                        Monitor.Wait(_mutex);
 89                    }
 90
 91                    if (_state >= StateClosing)
 92                    {
 93                        Debug.Assert(_exception is not null);
 94                        throw _exception;
 95                    }
 96                }
 97
 98                //
 99                // We start out in holding state.
 100                //
 101                setState(StateHolding);
 102            }
 103        }
 104        catch (LocalException ex)
 105        {
 106            exception(ex);
 107            waitUntilFinished();
 108            return;
 109        }
 110    }
 111
 112    internal void activate()
 113    {
 114        lock (_mutex)
 115        {
 116            if (_state <= StateNotValidated)
 117            {
 118                return;
 119            }
 120
 121            setState(StateActive);
 122        }
 123    }
 124
 125    internal void hold()
 126    {
 127        lock (_mutex)
 128        {
 129            if (_state <= StateNotValidated)
 130            {
 131                return;
 132            }
 133
 134            setState(StateHolding);
 135        }
 136    }
 137
 138    // DestructionReason.
 139    public const int ObjectAdapterDeactivated = 0;
 140    public const int CommunicatorDestroyed = 1;
 141
 142    internal void destroy(int reason)
 143    {
 144        lock (_mutex)
 145        {
 146            switch (reason)
 147            {
 148                case ObjectAdapterDeactivated:
 149                {
 150                    setState(StateClosing, new ObjectAdapterDeactivatedException(_adapter?.getName() ?? ""));
 151                    break;
 152                }
 153
 154                case CommunicatorDestroyed:
 155                {
 156                    setState(StateClosing, new CommunicatorDestroyedException());
 157                    break;
 158                }
 159            }
 160        }
 161    }
 162
 163    public void abort()
 164    {
 165        lock (_mutex)
 166        {
 167            setState(
 168                StateClosed,
 169                new ConnectionAbortedException(
 170                    "The connection was aborted by the application.",
 171                    closedByApplication: true));
 172        }
 173    }
 174
 175    public Task closeAsync()
 176    {
 177        lock (_mutex)
 178        {
 179            if (_state < StateClosing)
 180            {
 181                if (_asyncRequests.Count == 0)
 182                {
 183                    doApplicationClose();
 184                }
 185                else
 186                {
 187                    _closeRequested = true;
 188                    scheduleCloseTimer(); // we don't wait forever for outstanding invocations to complete
 189                }
 190            }
 191            // else nothing to do, already closing or closed.
 192        }
 193
 194        return _closed.Task;
 195    }
 196
 197    internal bool isActiveOrHolding()
 198    {
 199        lock (_mutex)
 200        {
 201            return _state > StateNotValidated && _state < StateClosing;
 202        }
 203    }
 204
 205    public void throwException()
 206    {
 207        lock (_mutex)
 208        {
 209            if (_exception is not null)
 210            {
 211                Debug.Assert(_state >= StateClosing);
 212                throw _exception;
 213            }
 214        }
 215    }
 216
 217    internal void waitUntilHolding()
 218    {
 219        lock (_mutex)
 220        {
 221            while (_state < StateHolding || _upcallCount > 0)
 222            {
 223                Monitor.Wait(_mutex);
 224            }
 225        }
 226    }
 227
 228    internal void waitUntilFinished()
 229    {
 230        lock (_mutex)
 231        {
 232            //
 233            // We wait indefinitely until the connection is finished and all
 234            // outstanding requests are completed. Otherwise we couldn't
 235            // guarantee that there are no outstanding calls when deactivate()
 236            // is called on the servant locators.
 237            //
 238            while (_state < StateFinished || _upcallCount > 0)
 239            {
 240                Monitor.Wait(_mutex);
 241            }
 242
 243            Debug.Assert(_state == StateFinished);
 244
 245            //
 246            // Clear the OA. See bug 1673 for the details of why this is necessary.
 247            //
 248            _adapter = null;
 249        }
 250    }
 251
 252    internal void updateObserver()
 253    {
 254        lock (_mutex)
 255        {
 256            if (_state < StateNotValidated || _state > StateClosed)
 257            {
 258                return;
 259            }
 260
 261            Debug.Assert(_instance.initializationData().observer is not null);
 262            _observer = _instance.initializationData().observer.getConnectionObserver(
 263                initConnectionInfo(),
 264                _endpoint,
 265                toConnectionState(_state),
 266                _observer);
 267            if (_observer is not null)
 268            {
 269                _observer.attach();
 270            }
 271            else
 272            {
 273                _writeStreamPos = -1;
 274                _readStreamPos = -1;
 275            }
 276        }
 277    }
 278
 279    internal int sendAsyncRequest(
 280        OutgoingAsyncBase og,
 281        bool compress,
 282        bool response,
 283        int batchRequestCount)
 284    {
 285        OutputStream os = og.getOs();
 286
 287        lock (_mutex)
 288        {
 289            //
 290            // If the connection is closed before we even have a chance
 291            // to send our request, we always try to send the request
 292            // again.
 293            //
 294            if (_exception is not null)
 295            {
 296                throw new RetryException(_exception);
 297            }
 298
 299            Debug.Assert(_state > StateNotValidated);
 300            Debug.Assert(_state < StateClosing);
 301
 302            //
 303            // Ensure the message isn't bigger than what we can send with the
 304            // transport.
 305            //
 306            _transceiver.checkSendSize(os.getBuffer());
 307
 308            //
 309            // Notify the request that it's cancelable with this connection.
 310            // This will throw if the request is canceled.
 311            //
 312            og.cancelable(this);
 313            int requestId = 0;
 314            if (response)
 315            {
 316                //
 317                // Create a new unique request ID.
 318                //
 319                requestId = _nextRequestId++;
 320                if (requestId <= 0)
 321                {
 322                    _nextRequestId = 1;
 323                    requestId = _nextRequestId++;
 324                }
 325
 326                //
 327                // Fill in the request ID.
 328                //
 329                os.pos(Protocol.headerSize);
 330                os.writeInt(requestId);
 331            }
 332            else if (batchRequestCount > 0)
 333            {
 334                os.pos(Protocol.headerSize);
 335                os.writeInt(batchRequestCount);
 336            }
 337
 338            og.attachRemoteObserver(initConnectionInfo(), _endpoint, requestId);
 339
 340            // We're just about to send a request, so we are not inactive anymore.
 341            cancelInactivityTimer();
 342
 343            int status = OutgoingAsyncBase.AsyncStatusQueued;
 344            try
 345            {
 346                var message = new OutgoingMessage(og, os, compress, requestId);
 347                status = sendMessage(message);
 348            }
 349            catch (LocalException ex)
 350            {
 351                setState(StateClosed, ex);
 352                Debug.Assert(_exception is not null);
 353                throw _exception;
 354            }
 355
 356            if (response)
 357            {
 358                //
 359                // Add to the async requests map.
 360                //
 361                _asyncRequests[requestId] = og;
 362            }
 363            return status;
 364        }
 365    }
 366
 367    internal BatchRequestQueue getBatchRequestQueue() => _batchRequestQueue;
 368
 369    public void flushBatchRequests(CompressBatch compress)
 370    {
 371        try
 372        {
 373            var completed = new FlushBatchTaskCompletionCallback();
 374            var outgoing = new ConnectionFlushBatchAsync(this, _instance, completed);
 375            outgoing.invoke(_flushBatchRequests_name, compress, true);
 376            completed.Task.Wait();
 377        }
 378        catch (AggregateException ex)
 379        {
 380            throw ex.InnerException;
 381        }
 382    }
 383
 384    public Task flushBatchRequestsAsync(
 385        CompressBatch compress,
 386        IProgress<bool> progress = null,
 387        CancellationToken cancel = default)
 388    {
 389        var completed = new FlushBatchTaskCompletionCallback(progress, cancel);
 390        var outgoing = new ConnectionFlushBatchAsync(this, _instance, completed);
 391        outgoing.invoke(_flushBatchRequests_name, compress, false);
 392        return completed.Task;
 393    }
 394
 395    private const string _flushBatchRequests_name = "flushBatchRequests";
 396
 397    public void disableInactivityCheck()
 398    {
 399        lock (_mutex)
 400        {
 401            cancelInactivityTimer();
 402            _inactivityTimeout = TimeSpan.Zero;
 403        }
 404    }
 405
 406    public void setCloseCallback(CloseCallback callback)
 407    {
 408        lock (_mutex)
 409        {
 410            if (_state >= StateClosed)
 411            {
 412                if (callback is not null)
 413                {
 414                    _threadPool.execute(
 415                        () =>
 416                        {
 417                            try
 418                            {
 419                                callback(this);
 420                            }
 421                            catch (System.Exception ex)
 422                            {
 423                                _logger.error("connection callback exception:\n" + ex + '\n' + _desc);
 424                            }
 425                        },
 426                        this);
 427                }
 428            }
 429            else
 430            {
 431                _closeCallback = callback;
 432            }
 433        }
 434    }
 435
 436    public void asyncRequestCanceled(OutgoingAsyncBase outAsync, LocalException ex)
 437    {
 438        //
 439        // NOTE: This isn't called from a thread pool thread.
 440        //
 441
 442        lock (_mutex)
 443        {
 444            if (_state >= StateClosed)
 445            {
 446                return; // The request has already been or will be shortly notified of the failure.
 447            }
 448
 449            OutgoingMessage o = _sendStreams.FirstOrDefault(m => m.outAsync == outAsync);
 450            if (o is not null)
 451            {
 452                if (o.requestId > 0)
 453                {
 454                    _asyncRequests.Remove(o.requestId);
 455                }
 456
 457                if (ex is ConnectionAbortedException)
 458                {
 459                    setState(StateClosed, ex);
 460                }
 461                else
 462                {
 463                    //
 464                    // If the request is being sent, don't remove it from the send streams,
 465                    // it will be removed once the sending is finished.
 466                    //
 467                    if (o == _sendStreams.First.Value)
 468                    {
 469                        o.canceled();
 470                    }
 471                    else
 472                    {
 473                        o.canceled();
 474                        _sendStreams.Remove(o);
 475                    }
 476                    if (outAsync.exception(ex))
 477                    {
 478                        outAsync.invokeExceptionAsync();
 479                    }
 480                }
 481
 482                if (_closeRequested && _state < StateClosing && _asyncRequests.Count == 0)
 483                {
 484                    doApplicationClose();
 485                }
 486                return;
 487            }
 488
 489            if (outAsync is OutgoingAsync)
 490            {
 491                foreach (KeyValuePair<int, OutgoingAsyncBase> kvp in _asyncRequests)
 492                {
 493                    if (kvp.Value == outAsync)
 494                    {
 495                        if (ex is ConnectionAbortedException)
 496                        {
 497                            setState(StateClosed, ex);
 498                        }
 499                        else
 500                        {
 501                            _asyncRequests.Remove(kvp.Key);
 502                            if (outAsync.exception(ex))
 503                            {
 504                                outAsync.invokeExceptionAsync();
 505                            }
 506                        }
 507
 508                        if (_closeRequested && _state < StateClosing && _asyncRequests.Count == 0)
 509                        {
 510                            doApplicationClose();
 511                        }
 512                        return;
 513                    }
 514                }
 515            }
 516        }
 517    }
 518
 519    internal EndpointI endpoint() => _endpoint; // No mutex protection necessary, _endpoint is immutable.
 520
 521    internal Connector connector() => _connector; // No mutex protection necessary, _endpoint is immutable.
 522
 523    public void setAdapter(ObjectAdapter adapter)
 524    {
 525        if (_connector is null) // server connection
 526        {
 527            throw new InvalidOperationException("setAdapter can only be called on a client connection");
 528        }
 529
 530        if (adapter is not null)
 531        {
 532            // Go through the adapter to set the adapter and servant manager on this connection
 533            // to ensure the object adapter is still active.
 534            adapter.setAdapterOnConnection(this);
 535        }
 536        else
 537        {
 538            lock (_mutex)
 539            {
 540                if (_state <= StateNotValidated || _state >= StateClosing)
 541                {
 542                    return;
 543                }
 544                _adapter = null;
 545            }
 546        }
 547
 548        //
 549        // We never change the thread pool with which we were initially
 550        // registered, even if we add or remove an object adapter.
 551        //
 552    }
 553
 554    public ObjectAdapter getAdapter()
 555    {
 556        lock (_mutex)
 557        {
 558            return _adapter;
 559        }
 560    }
 561
 562    public Endpoint getEndpoint() => _endpoint; // No mutex protection necessary, _endpoint is immutable.
 563
 564    public ObjectPrx createProxy(Identity id)
 565    {
 566        ObjectAdapter.checkIdentity(id);
 567        return new ObjectPrxHelper(_instance.referenceFactory().create(id, this));
 568    }
 569
 570    public void setAdapterFromAdapter(ObjectAdapter adapter)
 571    {
 572        lock (_mutex)
 573        {
 574            if (_state <= StateNotValidated || _state >= StateClosing)
 575            {
 576                return;
 577            }
 578            Debug.Assert(adapter is not null); // Called by ObjectAdapter::setAdapterOnConnection
 579            _adapter = adapter;
 580
 581            // Clear cached connection info (if any) as it's no longer accurate.
 582            _info = null;
 583        }
 584    }
 585
 586    //
 587    // Operations from EventHandler
 588    //
 589    public override bool startAsync(int operation, Ice.Internal.AsyncCallback completedCallback)
 590    {
 591        if (_state >= StateClosed)
 592        {
 593            return false;
 594        }
 595
 596        if (_threadHopRequired)
 597        {
 598            // Run the I/O on a .NET ThreadPool thread so it survives the initiating Ice worker exiting.
 599            // See ctor for when this is required.
 600            Task.Run(doIO);
 601        }
 602        else
 603        {
 604            doIO();
 605        }
 606
 607        return true;
 608
 609        void doIO()
 610        {
 611            lock (_mutex)
 612            {
 613                if (_state >= StateClosed)
 614                {
 615                    completedCallback(this);
 616                    return;
 617                }
 618
 619                try
 620                {
 621                    if ((operation & SocketOperation.Write) != 0)
 622                    {
 623                        if (_observer != null)
 624                        {
 625                            observerStartWrite(_writeStream.getBuffer());
 626                        }
 627
 628                        bool completedSynchronously =
 629                            _transceiver.startWrite(
 630                                _writeStream.getBuffer(),
 631                                completedCallback,
 632                                this,
 633                                out bool messageWritten);
 634
 635                        // If the startWrite call wrote the message, we assume the message is sent now for at-most-once
 636                        // semantics in the event the connection is closed while the message is still in _sendStreams.
 637                        if (messageWritten && _sendStreams.Count > 0)
 638                        {
 639                            // See finish() code.
 640                            _sendStreams.First.Value.isSent = true;
 641                        }
 642
 643                        if (completedSynchronously)
 644                        {
 645                            // If the write completed synchronously, we need to call the completedCallback.
 646                            completedCallback(this);
 647                        }
 648                    }
 649                    else if ((operation & SocketOperation.Read) != 0)
 650                    {
 651                        if (_observer != null && !_readHeader)
 652                        {
 653                            observerStartRead(_readStream.getBuffer());
 654                        }
 655
 656                        if (_transceiver.startRead(_readStream.getBuffer(), completedCallback, this))
 657                        {
 658                            completedCallback(this);
 659                        }
 660                    }
 661                }
 662                catch (LocalException ex)
 663                {
 664                    setState(StateClosed, ex);
 665                    completedCallback(this);
 666                }
 667            }
 668        }
 669    }
 670
 671    public override bool finishAsync(int operation)
 672    {
 673        if (_state >= StateClosed)
 674        {
 675            return false;
 676        }
 677
 678        try
 679        {
 680            if ((operation & SocketOperation.Write) != 0)
 681            {
 682                Ice.Internal.Buffer buf = _writeStream.getBuffer();
 683                int start = buf.b.position();
 684                _transceiver.finishWrite(buf);
 685                if (_instance.traceLevels().network >= 3 && buf.b.position() != start)
 686                {
 687                    var s = new StringBuilder("sent ");
 688                    s.Append(buf.b.position() - start);
 689                    if (!_endpoint.datagram())
 690                    {
 691                        s.Append(" of ");
 692                        s.Append(buf.b.limit() - start);
 693                    }
 694                    s.Append(" bytes via ");
 695                    s.Append(_endpoint.protocol());
 696                    s.Append('\n');
 697                    s.Append(ToString());
 698                    _instance.initializationData().logger.trace(_instance.traceLevels().networkCat, s.ToString());
 699                }
 700
 701                if (_observer is not null)
 702                {
 703                    observerFinishWrite(_writeStream.getBuffer());
 704                }
 705            }
 706            else if ((operation & SocketOperation.Read) != 0)
 707            {
 708                Ice.Internal.Buffer buf = _readStream.getBuffer();
 709                int start = buf.b.position();
 710                _transceiver.finishRead(buf);
 711                if (_instance.traceLevels().network >= 3 && buf.b.position() != start)
 712                {
 713                    var s = new StringBuilder("received ");
 714                    if (_endpoint.datagram())
 715                    {
 716                        s.Append(buf.b.limit());
 717                    }
 718                    else
 719                    {
 720                        s.Append(buf.b.position() - start);
 721                        s.Append(" of ");
 722                        s.Append(buf.b.limit() - start);
 723                    }
 724                    s.Append(" bytes via ");
 725                    s.Append(_endpoint.protocol());
 726                    s.Append('\n');
 727                    s.Append(ToString());
 728                    _instance.initializationData().logger.trace(_instance.traceLevels().networkCat, s.ToString());
 729                }
 730
 731                if (_observer is not null && !_readHeader)
 732                {
 733                    observerFinishRead(_readStream.getBuffer());
 734                }
 735            }
 736        }
 737        catch (LocalException ex)
 738        {
 739            setState(StateClosed, ex);
 740        }
 741        return _state < StateClosed;
 742    }
 743
 744    public override void message(ThreadPoolCurrent current)
 745    {
 746        StartCallback startCB = null;
 747        Queue<OutgoingMessage> sentCBs = null;
 748        var info = new MessageInfo();
 749        int upcallCount = 0;
 750
 751        using var msg = new ThreadPoolMessage(current, _mutex);
 752        lock (_mutex)
 753        {
 754            try
 755            {
 756                if (!msg.startIOScope())
 757                {
 758                    return;
 759                }
 760
 761                if (_state >= StateClosed)
 762                {
 763                    return;
 764                }
 765
 766                try
 767                {
 768                    int writeOp = SocketOperation.None;
 769                    int readOp = SocketOperation.None;
 770
 771                    // If writes are ready, write the data from the connection's write buffer (_writeStream)
 772                    if ((current.operation & SocketOperation.Write) != 0)
 773                    {
 774                        if (_observer is not null)
 775                        {
 776                            observerStartWrite(_writeStream.getBuffer());
 777                        }
 778                        writeOp = write(_writeStream.getBuffer());
 779                        if (_observer is not null && (writeOp & SocketOperation.Write) == 0)
 780                        {
 781                            observerFinishWrite(_writeStream.getBuffer());
 782                        }
 783                    }
 784
 785                    // If reads are ready, read the data into the connection's read buffer (_readStream). The data is
 786                    // read until:
 787                    // - the full message is read (the transport read returns SocketOperationNone) and
 788                    //   the read buffer is fully filled
 789                    // - the read operation on the transport can't continue without blocking
 790                    if ((current.operation & SocketOperation.Read) != 0)
 791                    {
 792                        while (true)
 793                        {
 794                            Ice.Internal.Buffer buf = _readStream.getBuffer();
 795
 796                            if (_observer is not null && !_readHeader)
 797                            {
 798                                observerStartRead(buf);
 799                            }
 800
 801                            readOp = read(buf);
 802                            if ((readOp & SocketOperation.Read) != 0)
 803                            {
 804                                // Can't continue without blocking, exit out of the loop.
 805                                break;
 806                            }
 807                            if (_observer is not null && !_readHeader)
 808                            {
 809                                Debug.Assert(!buf.b.hasRemaining());
 810                                observerFinishRead(buf);
 811                            }
 812
 813                            // If read header is true, we're reading a new Ice protocol message and we need to read
 814                            // the message header.
 815                            if (_readHeader)
 816                            {
 817                                // The next read will read the remainder of the message.
 818                                _readHeader = false;
 819
 820                                _observer?.receivedBytes(Protocol.headerSize);
 821
 822                                //
 823                                // Connection is validated on first message. This is only used by
 824                                // setState() to check whether or not we can print a connection
 825                                // warning (a client might close the connection forcefully if the
 826                                // connection isn't validated, we don't want to print a warning
 827                                // in this case).
 828                                //
 829                                _validated = true;
 830
 831                                // Full header should be read because the size of _readStream is always headerSize (14)
 832                                // when reading a new message (see the code that sets _readHeader = true).
 833                                int pos = _readStream.pos();
 834                                if (pos < Protocol.headerSize)
 835                                {
 836                                    //
 837                                    // This situation is possible for small UDP packets.
 838                                    //
 839                                    throw new MarshalException("Received Ice message with too few bytes in header.");
 840                                }
 841
 842                                // Decode the header.
 843                                _readStream.pos(0);
 844                                byte[] m = new byte[4];
 845                                m[0] = _readStream.readByte();
 846                                m[1] = _readStream.readByte();
 847                                m[2] = _readStream.readByte();
 848                                m[3] = _readStream.readByte();
 849                                if (m[0] != Protocol.magic[0] || m[1] != Protocol.magic[1] ||
 850                                m[2] != Protocol.magic[2] || m[3] != Protocol.magic[3])
 851                                {
 852                                    throw new ProtocolException(
 853                                        $"Bad magic in message header: {m[0]:X2} {m[1]:X2} {m[2]:X2} {m[3]:X2}");
 854                                }
 855
 856                                var pv = new ProtocolVersion(_readStream);
 857                                if (pv != Protocol.currentProtocol)
 858                                {
 859                                    throw new MarshalException(
 860                                        $"Invalid protocol version in message header: {pv.major}.{pv.minor}");
 861                                }
 862                                var ev = new EncodingVersion(_readStream);
 863                                if (ev != Protocol.currentProtocolEncoding)
 864                                {
 865                                    throw new MarshalException(
 866                                        $"Invalid protocol encoding version in message header: {ev.major}.{ev.minor}");
 867                                }
 868
 869                                _readStream.readByte(); // messageType
 870                                _readStream.readByte(); // compress
 871                                int size = _readStream.readInt();
 872                                if (size < Protocol.headerSize)
 873                                {
 874                                    throw new MarshalException($"Received Ice message with unexpected size {size}.");
 875                                }
 876
 877                                // Resize the read buffer to the message size.
 878                                if (size > _messageSizeMax)
 879                                {
 880                                    Ex.throwMemoryLimitException(size, _messageSizeMax);
 881                                }
 882                                if (size > _readStream.size())
 883                                {
 884                                    _readStream.resize(size);
 885                                }
 886                                _readStream.pos(pos);
 887                            }
 888
 889                            if (buf.b.hasRemaining())
 890                            {
 891                                if (_endpoint.datagram())
 892                                {
 893                                    throw new DatagramLimitException(); // The message was truncated.
 894                                }
 895                                continue;
 896                            }
 897                            break;
 898                        }
 899                    }
 900
 901                    // readOp and writeOp are set to the operations that the transport read or write calls from above
 902                    // returned. They indicate which operations will need to be monitored by the thread pool's selector
 903                    // when this method returns.
 904                    int newOp = readOp | writeOp;
 905
 906                    // Operations that are ready. For example, if message was called with SocketOperationRead and the
 907                    // transport read returned SocketOperationNone, reads are considered done: there's no additional
 908                    // data to read.
 909                    int readyOp = current.operation & ~newOp;
 910
 911                    if (_state <= StateNotValidated)
 912                    {
 913                        // If the connection is still not validated and there's still data to read or write, continue
 914                        // waiting for data to read or write.
 915                        if (newOp != 0)
 916                        {
 917                            _threadPool.update(this, current.operation, newOp);
 918                            return;
 919                        }
 920
 921                        // Initialize the connection if it's not initialized yet.
 922                        if (_state == StateNotInitialized && !initialize(current.operation))
 923                        {
 924                            return;
 925                        }
 926
 927                        // Validate the connection if it's not validated yet.
 928                        if (_state <= StateNotValidated && !validate(current.operation))
 929                        {
 930                            return;
 931                        }
 932
 933                        // The connection is validated and doesn't need additional data to be read or written. So
 934                        // unregister it from the thread pool's selector.
 935                        _threadPool.unregister(this, current.operation);
 936
 937                        //
 938                        // We start out in holding state.
 939                        //
 940                        setState(StateHolding);
 941                        if (_startCallback is not null)
 942                        {
 943                            startCB = _startCallback;
 944                            _startCallback = null;
 945                            ++upcallCount;
 946                        }
 947                    }
 948                    else
 949                    {
 950                        Debug.Assert(_state <= StateClosingPending);
 951
 952                        //
 953                        // We parse messages first, if we receive a close
 954                        // connection message we won't send more messages.
 955                        //
 956                        if ((readyOp & SocketOperation.Read) != 0)
 957                        {
 958                            // At this point, the protocol message is fully read and can therefore be decoded by
 959                            // parseMessage. parseMessage returns the operation to wait for readiness next.
 960                            newOp |= parseMessage(ref info);
 961                            upcallCount += info.upcallCount;
 962                        }
 963
 964                        if ((readyOp & SocketOperation.Write) != 0)
 965                        {
 966                            // At this point the message from _writeStream is fully written and the next message can be
 967                            // written.
 968
 969                            newOp |= sendNextMessage(out sentCBs);
 970                            if (sentCBs is not null)
 971                            {
 972                                ++upcallCount;
 973                            }
 974                        }
 975
 976                        // If the connection is not closed yet, we can update the thread pool selector to wait for
 977                        // readiness of read, write or both operations.
 978                        if (_state < StateClosed)
 979                        {
 980                            _threadPool.update(this, current.operation, newOp);
 981                        }
 982                    }
 983
 984                    if (upcallCount == 0)
 985                    {
 986                        return; // Nothing to execute, we're done!
 987                    }
 988
 989                    _upcallCount += upcallCount;
 990
 991                    // There's something to execute so we mark IO as completed to elect a new leader thread and let IO
 992                    // be performed on this new leader thread while this thread continues with executing the upcalls.
 993                    msg.ioCompleted();
 994                }
 995                catch (DatagramLimitException) // Expected.
 996                {
 997                    if (_warnUdp)
 998                    {
 999                        _logger.warning($"maximum datagram size of {_readStream.pos()} exceeded");
 1000                    }
 1001                    _readStream.resize(Protocol.headerSize);
 1002                    _readStream.pos(0);
 1003                    _readHeader = true;
 1004                    return;
 1005                }
 1006                catch (SocketException ex)
 1007                {
 1008                    setState(StateClosed, ex);
 1009                    return;
 1010                }
 1011                catch (LocalException ex)
 1012                {
 1013                    if (_endpoint.datagram())
 1014                    {
 1015                        if (_warn)
 1016                        {
 1017                            _logger.warning($"datagram connection exception:\n{ex}\n{_desc}");
 1018                        }
 1019                        _readStream.resize(Protocol.headerSize);
 1020                        _readStream.pos(0);
 1021                        _readHeader = true;
 1022                    }
 1023                    else
 1024                    {
 1025                        setState(StateClosed, ex);
 1026                    }
 1027                    return;
 1028                }
 1029            }
 1030            finally
 1031            {
 1032                msg.finishIOScope();
 1033            }
 1034        }
 1035
 1036        _threadPool.executeFromThisThread(() => upcall(startCB, sentCBs, info), this);
 1037    }
 1038
 1039    private void upcall(StartCallback startCB, Queue<OutgoingMessage> sentCBs, MessageInfo info)
 1040    {
 1041        int completedUpcallCount = 0;
 1042
 1043        //
 1044        // Notify the factory that the connection establishment and
 1045        // validation has completed.
 1046        //
 1047        if (startCB is not null)
 1048        {
 1049            startCB.connectionStartCompleted(this);
 1050            ++completedUpcallCount;
 1051        }
 1052
 1053        //
 1054        // Notify AMI calls that the message was sent.
 1055        //
 1056        if (sentCBs is not null)
 1057        {
 1058            foreach (OutgoingMessage m in sentCBs)
 1059            {
 1060                if (m.invokeSent)
 1061                {
 1062                    m.outAsync.invokeSent();
 1063                }
 1064                if (m.receivedReply)
 1065                {
 1066                    var outAsync = (OutgoingAsync)m.outAsync;
 1067                    if (outAsync.response())
 1068                    {
 1069                        outAsync.invokeResponse();
 1070                    }
 1071                }
 1072            }
 1073            ++completedUpcallCount;
 1074        }
 1075
 1076        //
 1077        // Asynchronous replies must be handled outside the thread
 1078        // synchronization, so that nested calls are possible.
 1079        //
 1080        if (info.outAsync is not null)
 1081        {
 1082            info.outAsync.invokeResponse();
 1083            ++completedUpcallCount;
 1084        }
 1085
 1086        //
 1087        // Method invocation (or multiple invocations for batch messages)
 1088        // must be done outside the thread synchronization, so that nested
 1089        // calls are possible.
 1090        //
 1091        if (info.requestCount > 0)
 1092        {
 1093            dispatchAll(info.stream, info.requestCount, info.requestId, info.compress, info.adapter);
 1094        }
 1095
 1096        //
 1097        // Decrease the upcall count.
 1098        //
 1099        bool finished = false;
 1100        if (completedUpcallCount > 0)
 1101        {
 1102            lock (_mutex)
 1103            {
 1104                _upcallCount -= completedUpcallCount;
 1105                if (_upcallCount == 0)
 1106                {
 1107                    // Only initiate shutdown if not already initiated. It might have already been initiated if the sent
 1108                    // callback or AMI callback was called when the connection was in the closing state.
 1109                    if (_state == StateClosing)
 1110                    {
 1111                        try
 1112                        {
 1113                            initiateShutdown();
 1114                        }
 1115                        catch (Ice.LocalException ex)
 1116                        {
 1117                            setState(StateClosed, ex);
 1118                        }
 1119                    }
 1120                    else if (_state == StateFinished)
 1121                    {
 1122                        finished = true;
 1123                        _observer?.detach();
 1124                    }
 1125                    Monitor.PulseAll(_mutex);
 1126                }
 1127            }
 1128        }
 1129
 1130        if (finished && _removeFromFactory is not null)
 1131        {
 1132            _removeFromFactory(this);
 1133        }
 1134    }
 1135
 1136    public override void finished(ThreadPoolCurrent current)
 1137    {
 1138        // Lock the connection here to ensure setState() completes before the code below is executed. This method can
 1139        // be called by the thread pool as soon as setState() calls _threadPool->finish(...). There's no need to lock
 1140        // the mutex for the remainder of the code because the data members accessed by finish() are immutable once
 1141        // _state == StateClosed (and we don't want to hold the mutex when calling upcalls).
 1142        lock (_mutex)
 1143        {
 1144            Debug.Assert(_state == StateClosed);
 1145        }
 1146
 1147        //
 1148        // If there are no callbacks to call, we don't call ioCompleted() since we're not going
 1149        // to call code that will potentially block (this avoids promoting a new leader and
 1150        // unnecessary thread creation, especially if this is called on shutdown).
 1151        //
 1152        if (_startCallback is null && _sendStreams.Count == 0 && _asyncRequests.Count == 0 && _closeCallback is null)
 1153        {
 1154            finish();
 1155            return;
 1156        }
 1157
 1158        current.ioCompleted();
 1159        _threadPool.executeFromThisThread(finish, this);
 1160    }
 1161
 1162    private void finish()
 1163    {
 1164        if (!_initialized)
 1165        {
 1166            if (_instance.traceLevels().network >= 2)
 1167            {
 1168                var s = new StringBuilder("failed to ");
 1169                s.Append(_connector is not null ? "establish" : "accept");
 1170                s.Append(' ');
 1171                s.Append(_endpoint.protocol());
 1172                s.Append(" connection\n");
 1173                s.Append(ToString());
 1174                s.Append('\n');
 1175                s.Append(_exception);
 1176                _instance.initializationData().logger.trace(_instance.traceLevels().networkCat, s.ToString());
 1177            }
 1178        }
 1179        else
 1180        {
 1181            if (_instance.traceLevels().network >= 1)
 1182            {
 1183                var s = new StringBuilder("closed ");
 1184                s.Append(_endpoint.protocol());
 1185                s.Append(" connection\n");
 1186                s.Append(ToString());
 1187
 1188                // Trace the cause of most connection closures.
 1189                if (!(_exception is CommunicatorDestroyedException || _exception is ObjectAdapterDeactivatedException))
 1190                {
 1191                    s.Append('\n');
 1192                    s.Append(_exception);
 1193                }
 1194
 1195                _instance.initializationData().logger.trace(_instance.traceLevels().networkCat, s.ToString());
 1196            }
 1197        }
 1198
 1199        _startCallback?.connectionStartFailed(this, _exception);
 1200        _startCallback = null;
 1201
 1202        if (_sendStreams.Count > 0)
 1203        {
 1204            if (!_writeStream.isEmpty())
 1205            {
 1206                //
 1207                // Return the stream to the outgoing call. This is important for
 1208                // retriable AMI calls which are not marshaled again.
 1209                //
 1210                OutgoingMessage message = _sendStreams.First.Value;
 1211                _writeStream.swap(message.stream);
 1212
 1213                //
 1214                // The current message might be sent but not yet removed from _sendStreams. If
 1215                // the response has been received in the meantime, we remove the message from
 1216                // _sendStreams to not call finished on a message which is already done.
 1217                //
 1218                if (message.isSent || message.receivedReply)
 1219                {
 1220                    if (message.sent() && message.invokeSent)
 1221                    {
 1222                        message.outAsync.invokeSent();
 1223                    }
 1224                    if (message.receivedReply)
 1225                    {
 1226                        var outAsync = (OutgoingAsync)message.outAsync;
 1227                        if (outAsync.response())
 1228                        {
 1229                            outAsync.invokeResponse();
 1230                        }
 1231                    }
 1232                    _sendStreams.RemoveFirst();
 1233                }
 1234            }
 1235
 1236            foreach (OutgoingMessage o in _sendStreams)
 1237            {
 1238                o.completed(_exception);
 1239                if (o.requestId > 0) // Make sure finished isn't called twice.
 1240                {
 1241                    _asyncRequests.Remove(o.requestId);
 1242                }
 1243            }
 1244            _sendStreams.Clear(); // Must be cleared before _requests because of Outgoing* references in OutgoingMessage
 1245        }
 1246
 1247        foreach (OutgoingAsyncBase o in _asyncRequests.Values)
 1248        {
 1249            if (o.exception(_exception))
 1250            {
 1251                o.invokeException();
 1252            }
 1253        }
 1254        _asyncRequests.Clear();
 1255
 1256        //
 1257        // Don't wait to be reaped to reclaim memory allocated by read/write streams.
 1258        //
 1259        _writeStream.clear();
 1260        _writeStream.getBuffer().clear();
 1261        _readStream.clear();
 1262        _readStream.getBuffer().clear();
 1263
 1264        if (_exception is ConnectionClosedException or
 1265            CloseConnectionException or
 1266            CommunicatorDestroyedException or
 1267            ObjectAdapterDeactivatedException)
 1268        {
 1269            // Can execute synchronously. Note that we're not within a lock(this) here.
 1270            _closed.SetResult();
 1271        }
 1272        else
 1273        {
 1274            Debug.Assert(_exception is not null);
 1275            _closed.SetException(_exception);
 1276        }
 1277
 1278        if (_closeCallback is not null)
 1279        {
 1280            try
 1281            {
 1282                _closeCallback(this);
 1283            }
 1284            catch (System.Exception ex)
 1285            {
 1286                _logger.error("connection callback exception:\n" + ex + '\n' + _desc);
 1287            }
 1288            _closeCallback = null;
 1289        }
 1290
 1291        //
 1292        // This must be done last as this will cause waitUntilFinished() to return (and communicator
 1293        // objects such as the timer might be destroyed too).
 1294        //
 1295        bool finished = false;
 1296        lock (_mutex)
 1297        {
 1298            setState(StateFinished);
 1299
 1300            if (_upcallCount == 0)
 1301            {
 1302                finished = true;
 1303                _observer?.detach();
 1304            }
 1305        }
 1306
 1307        if (finished && _removeFromFactory is not null)
 1308        {
 1309            _removeFromFactory(this);
 1310        }
 1311    }
 1312
 1313    /// <inheritdoc/>
 1314    public override string ToString() => _desc; // No mutex lock, _desc is immutable.
 1315
 1316    /// <inheritdoc/>
 1317    public string type() => _type; // No mutex lock, _type is immutable.
 1318
 1319    /// <inheritdoc/>
 1320    public ConnectionInfo getInfo()
 1321    {
 1322        lock (_mutex)
 1323        {
 1324            if (_state >= StateClosed)
 1325            {
 1326                throw _exception;
 1327            }
 1328            return initConnectionInfo();
 1329        }
 1330    }
 1331
 1332    /// <inheritdoc/>
 1333    public void setBufferSize(int rcvSize, int sndSize)
 1334    {
 1335        lock (_mutex)
 1336        {
 1337            if (_state >= StateClosed)
 1338            {
 1339                throw _exception;
 1340            }
 1341            try
 1342            {
 1343                _transceiver.setBufferSize(rcvSize, sndSize);
 1344            }
 1345            catch (LocalException ex)
 1346            {
 1347                // The failing call may have closed the socket, so close the connection as well.
 1348                setState(StateClosed, ex);
 1349                throw;
 1350            }
 1351            _info = null; // Invalidate the cached connection info
 1352        }
 1353    }
 1354
 1355    public void exception(LocalException ex)
 1356    {
 1357        lock (_mutex)
 1358        {
 1359            setState(StateClosed, ex);
 1360        }
 1361    }
 1362
 1363    public Ice.Internal.ThreadPool getThreadPool() => _threadPool;
 1364
 1365    internal ConnectionI(
 1366        Instance instance,
 1367        Transceiver transceiver,
 1368        Connector connector, // null for incoming connections, non-null for outgoing connections
 1369        EndpointI endpoint,
 1370        ObjectAdapter adapter,
 1371        Action<ConnectionI> removeFromFactory, // can be null
 1372        ConnectionOptions options)
 1373    {
 1374        _instance = instance;
 1375        _desc = transceiver.ToString();
 1376        _type = transceiver.protocol();
 1377        _connector = connector;
 1378        _endpoint = endpoint;
 1379        _adapter = adapter;
 1380        InitializationData initData = instance.initializationData();
 1381        _logger = initData.logger; // Cached for better performance.
 1382        _traceLevels = instance.traceLevels(); // Cached for better performance.
 1383        _connectTimeout = options.connectTimeout;
 1384        _closeTimeout = options.closeTimeout; // not used for datagram connections
 1385        // suppress inactivity timeout for datagram connections
 1386        _inactivityTimeout = endpoint.datagram() ? TimeSpan.Zero : options.inactivityTimeout;
 1387        _maxDispatches = options.maxDispatches;
 1388        _removeFromFactory = removeFromFactory;
 1389        _warn = initData.properties.getIcePropertyAsInt("Ice.Warn.Connections") > 0;
 1390        _warnUdp = initData.properties.getIcePropertyAsInt("Ice.Warn.Datagrams") > 0;
 1391        _nextRequestId = 1;
 1392        _messageSizeMax = connector is null ? adapter.messageSizeMax() : instance.messageSizeMax();
 1393        _batchRequestQueue = new BatchRequestQueue(instance, _endpoint.datagram());
 1394        _readStream = new InputStream(instance, Protocol.currentProtocolEncoding);
 1395        _readHeader = false;
 1396        _readStreamPos = -1;
 1397        _writeStream = new OutputStream(); // temporary stream
 1398        _writeStreamPos = -1;
 1399        _upcallCount = 0;
 1400        _state = StateNotInitialized;
 1401
 1402        _compressionLevel = initData.properties.getIcePropertyAsInt("Ice.Compression.Level");
 1403        if (_compressionLevel < 1)
 1404        {
 1405            _compressionLevel = 1;
 1406        }
 1407        else if (_compressionLevel > 9)
 1408        {
 1409            _compressionLevel = 9;
 1410        }
 1411
 1412        if (options.idleTimeout > TimeSpan.Zero && !endpoint.datagram())
 1413        {
 1414            _idleTimeoutTransceiver = new IdleTimeoutTransceiverDecorator(
 1415                transceiver,
 1416                this,
 1417                options.idleTimeout,
 1418                options.enableIdleCheck);
 1419            transceiver = _idleTimeoutTransceiver;
 1420        }
 1421        _transceiver = transceiver;
 1422
 1423        if (connector is null)
 1424        {
 1425            // adapter is always set for incoming connections
 1426            Debug.Assert(adapter is not null);
 1427            _threadPool = adapter.getThreadPool();
 1428        }
 1429        else
 1430        {
 1431            // we use the client thread pool for outgoing connections, even if there is an
 1432            // object adapter with its own thread pool.
 1433            _threadPool = instance.clientThreadPool();
 1434        }
 1435        // On Windows, async socket I/O initiated on a thread that subsequently terminates is cancelled by the OS
 1436        // with SocketError.OperationAborted. The Ice thread pool can reap workers idle past ThreadIdleTime when
 1437        // SizeMax > 1, so in that combination we hop the I/O onto the .NET ThreadPool (whose threads are managed
 1438        // by the runtime and not reaped while owning pending I/O). Other platforms and fixed-size Ice pools don't
 1439        // need the hop. See startAsync.
 1440        _threadHopRequired = AssemblyUtil.isWindows && _threadPool.canShrink;
 1441
 1442        // initialize only resets the handler state; it doesn't throw.
 1443        _threadPool.initialize(this);
 1444    }
 1445
 1446    /// <summary>
 1447    /// Aborts the connection with a <see cref="ConnectionAbortedException" /> if the connection is active and
 1448    /// does not receive a byte for some time. See the IdleTimeoutTransceiverDecorator.
 1449    /// </summary>
 1450    internal void idleCheck(TimeSpan idleTimeout)
 1451    {
 1452        lock (_mutex)
 1453        {
 1454            if (_state == StateActive && _idleTimeoutTransceiver!.idleCheckEnabled)
 1455            {
 1456                int idleTimeoutInSeconds = (int)idleTimeout.TotalSeconds;
 1457
 1458                setState(
 1459                    StateClosed,
 1460                    new ConnectionAbortedException(
 1461                        $"Connection aborted by the idle check because it did not receive any bytes for {idleTimeoutInSe
 1462                        closedByApplication: false));
 1463            }
 1464            // else nothing to do
 1465        }
 1466    }
 1467
 1468    internal void sendHeartbeat()
 1469    {
 1470        Debug.Assert(!_endpoint.datagram());
 1471
 1472        lock (_mutex)
 1473        {
 1474            if (_state == StateActive || _state == StateHolding || _state == StateClosing)
 1475            {
 1476                // We check if the connection has become inactive.
 1477                if (
 1478                    _inactivityTimer is null &&           // timer not already scheduled
 1479                    _inactivityTimeout > TimeSpan.Zero && // inactivity timeout is enabled
 1480                    _state == StateActive &&              // only schedule the timer if the connection is active
 1481                    _dispatchCount == 0 &&                // no pending dispatch
 1482                    _asyncRequests.Count == 0 &&          // no pending invocation
 1483                    _readHeader &&                        // we're not waiting for the remainder of an incoming message
 1484                    _sendStreams.Count <= 1)              // there is at most one pending outgoing message
 1485                {
 1486                    // We may become inactive while the peer is back-pressuring us. In this case, we only schedule the
 1487                    // inactivity timer if there is no pending outgoing message or the pending outgoing message is a
 1488                    // heartbeat.
 1489
 1490                    // The stream of the first _sendStreams message is in _writeStream.
 1491                    if (_sendStreams.Count == 0 || isHeartbeat(_writeStream))
 1492                    {
 1493                        scheduleInactivityTimer();
 1494                    }
 1495                }
 1496
 1497                // We send a heartbeat to the peer to generate a "write" on the connection. This write in turns creates
 1498                // a read on the peer, and resets the peer's idle check timer. When _sendStreams is not empty, there is
 1499                // already an outstanding write, so we don't need to send a heartbeat. It's possible the first message
 1500                // of _sendStreams was already sent but not yet removed from _sendStreams: it means the last write
 1501                // occurred very recently, which is good enough with respect to the idle check.
 1502                // As a result of this optimization, the only possible heartbeat in _sendStreams is the first
 1503                // _sendStreams message.
 1504                if (_sendStreams.Count == 0)
 1505                {
 1506                    var os = new OutputStream(Protocol.currentProtocolEncoding);
 1507                    os.writeBlob(Protocol.magic);
 1508                    ProtocolVersion.ice_write(os, Protocol.currentProtocol);
 1509                    EncodingVersion.ice_write(os, Protocol.currentProtocolEncoding);
 1510                    os.writeByte(Protocol.validateConnectionMsg);
 1511                    os.writeByte(0);
 1512                    os.writeInt(Protocol.headerSize); // Message size.
 1513                    try
 1514                    {
 1515                        _ = sendMessage(new OutgoingMessage(os, compress: false));
 1516                    }
 1517                    catch (LocalException ex)
 1518                    {
 1519                        setState(StateClosed, ex);
 1520                    }
 1521                }
 1522            }
 1523            // else nothing to do
 1524        }
 1525
 1526        static bool isHeartbeat(OutputStream stream) =>
 1527            stream.getBuffer().b.get(8) == Protocol.validateConnectionMsg;
 1528    }
 1529
 1530    private const int StateNotInitialized = 0;
 1531    private const int StateNotValidated = 1;
 1532    private const int StateActive = 2;
 1533    private const int StateHolding = 3;
 1534    private const int StateClosing = 4;
 1535    private const int StateClosingPending = 5;
 1536    private const int StateClosed = 6;
 1537    private const int StateFinished = 7;
 1538
 1539    private static ConnectionState toConnectionState(int state) => connectionStateMap[state];
 1540
 1541    private void setState(int state, LocalException ex)
 1542    {
 1543        //
 1544        // If setState() is called with an exception, then only closed
 1545        // and closing states are permissible.
 1546        //
 1547        Debug.Assert(state >= StateClosing);
 1548
 1549        if (_state == state) // Don't switch twice.
 1550        {
 1551            return;
 1552        }
 1553
 1554        if (_exception is null)
 1555        {
 1556            //
 1557            // If we are in closed state, an exception must be set.
 1558            //
 1559            Debug.Assert(_state != StateClosed);
 1560
 1561            _exception = ex;
 1562
 1563            //
 1564            // We don't warn if we are not validated.
 1565            //
 1566            if (_warn && _validated)
 1567            {
 1568                //
 1569                // Don't warn about certain expected exceptions.
 1570                //
 1571                if (!(_exception is CloseConnectionException ||
 1572                     _exception is ConnectionClosedException ||
 1573                     _exception is CommunicatorDestroyedException ||
 1574                     _exception is ObjectAdapterDeactivatedException ||
 1575                     (_exception is ConnectionLostException && _state >= StateClosing)))
 1576                {
 1577                    warning("connection exception", _exception);
 1578                }
 1579            }
 1580        }
 1581
 1582        //
 1583        // We must set the new state before we notify requests of any
 1584        // exceptions. Otherwise new requests may retry on a
 1585        // connection that is not yet marked as closed or closing.
 1586        //
 1587        setState(state);
 1588    }
 1589
 1590    private void setState(int state)
 1591    {
 1592        //
 1593        // We don't want to send close connection messages if the endpoint
 1594        // only supports oneway transmission from client to server.
 1595        //
 1596        if (_endpoint.datagram() && state == StateClosing)
 1597        {
 1598            state = StateClosed;
 1599        }
 1600
 1601        //
 1602        // Skip graceful shutdown if we are destroyed before validation.
 1603        //
 1604        if (_state <= StateNotValidated && state == StateClosing)
 1605        {
 1606            state = StateClosed;
 1607        }
 1608
 1609        if (_state == state) // Don't switch twice.
 1610        {
 1611            return;
 1612        }
 1613
 1614        if (state > StateActive)
 1615        {
 1616            // Dispose the inactivity timer, if not null.
 1617            cancelInactivityTimer();
 1618        }
 1619
 1620        try
 1621        {
 1622            switch (state)
 1623            {
 1624                case StateNotInitialized:
 1625                {
 1626                    Debug.Assert(false);
 1627                    break;
 1628                }
 1629
 1630                case StateNotValidated:
 1631                {
 1632                    if (_state != StateNotInitialized)
 1633                    {
 1634                        Debug.Assert(_state == StateClosed);
 1635                        return;
 1636                    }
 1637                    break;
 1638                }
 1639
 1640                case StateActive:
 1641                {
 1642                    //
 1643                    // Can only switch to active from holding or not validated.
 1644                    //
 1645                    if (_state != StateHolding && _state != StateNotValidated)
 1646                    {
 1647                        return;
 1648                    }
 1649
 1650                    if (_maxDispatches <= 0 || _dispatchCount < _maxDispatches)
 1651                    {
 1652                        _threadPool.register(this, SocketOperation.Read);
 1653                        _idleTimeoutTransceiver?.enableIdleCheck();
 1654                    }
 1655                    // else don't resume reading since we're at or over the _maxDispatches limit.
 1656
 1657                    break;
 1658                }
 1659
 1660                case StateHolding:
 1661                {
 1662                    //
 1663                    // Can only switch to holding from active or not validated.
 1664                    //
 1665                    if (_state != StateActive && _state != StateNotValidated)
 1666                    {
 1667                        return;
 1668                    }
 1669
 1670                    if (_state == StateActive && (_maxDispatches <= 0 || _dispatchCount < _maxDispatches))
 1671                    {
 1672                        _threadPool.unregister(this, SocketOperation.Read);
 1673                        _idleTimeoutTransceiver?.disableIdleCheck();
 1674                    }
 1675                    // else reads are already disabled because the _maxDispatches limit is reached or exceeded.
 1676
 1677                    break;
 1678                }
 1679
 1680                case StateClosing:
 1681                case StateClosingPending:
 1682                {
 1683                    //
 1684                    // Can't change back from closing pending.
 1685                    //
 1686                    if (_state >= StateClosingPending)
 1687                    {
 1688                        return;
 1689                    }
 1690                    break;
 1691                }
 1692
 1693                case StateClosed:
 1694                {
 1695                    if (_state == StateFinished)
 1696                    {
 1697                        return;
 1698                    }
 1699
 1700                    _batchRequestQueue.destroy(_exception);
 1701                    _threadPool.finish(this);
 1702                    _transceiver.close();
 1703                    break;
 1704                }
 1705
 1706                case StateFinished:
 1707                {
 1708                    Debug.Assert(_state == StateClosed);
 1709                    _transceiver.destroy();
 1710                    break;
 1711                }
 1712            }
 1713        }
 1714        catch (LocalException ex)
 1715        {
 1716            _logger.error("unexpected connection exception:\n" + ex + "\n" + _desc);
 1717        }
 1718
 1719        if (_instance.initializationData().observer is not null)
 1720        {
 1721            ConnectionState oldState = toConnectionState(_state);
 1722            ConnectionState newState = toConnectionState(state);
 1723            if (oldState != newState)
 1724            {
 1725                _observer = _instance.initializationData().observer.getConnectionObserver(
 1726                    initConnectionInfo(),
 1727                    _endpoint,
 1728                    newState,
 1729                    _observer);
 1730                if (_observer is not null)
 1731                {
 1732                    _observer.attach();
 1733                }
 1734                else
 1735                {
 1736                    _writeStreamPos = -1;
 1737                    _readStreamPos = -1;
 1738                }
 1739            }
 1740            if (_observer is not null && state == StateClosed && _exception is not null)
 1741            {
 1742                if (!(_exception is CloseConnectionException ||
 1743                     _exception is ConnectionClosedException ||
 1744                     _exception is CommunicatorDestroyedException ||
 1745                     _exception is ObjectAdapterDeactivatedException ||
 1746                     (_exception is ConnectionLostException && _state >= StateClosing)))
 1747                {
 1748                    _observer.failed(_exception.ice_id());
 1749                }
 1750            }
 1751        }
 1752        _state = state;
 1753
 1754        Monitor.PulseAll(_mutex);
 1755
 1756        if (_state == StateClosing && _upcallCount == 0)
 1757        {
 1758            try
 1759            {
 1760                initiateShutdown();
 1761            }
 1762            catch (LocalException ex)
 1763            {
 1764                setState(StateClosed, ex);
 1765            }
 1766        }
 1767    }
 1768
 1769    private void initiateShutdown()
 1770    {
 1771        Debug.Assert(_state == StateClosing && _upcallCount == 0);
 1772
 1773        if (_shutdownInitiated)
 1774        {
 1775            return;
 1776        }
 1777        _shutdownInitiated = true;
 1778
 1779        if (!_endpoint.datagram())
 1780        {
 1781            //
 1782            // Before we shut down, we send a close connection message.
 1783            //
 1784            var os = new OutputStream(Protocol.currentProtocolEncoding);
 1785            os.writeBlob(Protocol.magic);
 1786            ProtocolVersion.ice_write(os, Protocol.currentProtocol);
 1787            EncodingVersion.ice_write(os, Protocol.currentProtocolEncoding);
 1788            os.writeByte(Protocol.closeConnectionMsg);
 1789            os.writeByte(0); // Compression status: always zero for close connection.
 1790            os.writeInt(Protocol.headerSize); // Message size.
 1791
 1792            scheduleCloseTimer();
 1793
 1794            if ((sendMessage(new OutgoingMessage(os, compress: false)) & OutgoingAsyncBase.AsyncStatusSent) != 0)
 1795            {
 1796                setState(StateClosingPending);
 1797
 1798                //
 1799                // Notify the transceiver of the graceful connection closure.
 1800                //
 1801                int op = _transceiver.closing(true, _exception);
 1802                if (op != 0)
 1803                {
 1804                    _threadPool.register(this, op);
 1805                }
 1806            }
 1807        }
 1808    }
 1809
 1810    private bool initialize(int operation)
 1811    {
 1812        int s = _transceiver.initialize(_readStream.getBuffer(), _writeStream.getBuffer(), ref _hasMoreData);
 1813        if (s != SocketOperation.None)
 1814        {
 1815            _threadPool.update(this, operation, s);
 1816            return false;
 1817        }
 1818
 1819        //
 1820        // Update the connection description once the transceiver is initialized.
 1821        //
 1822        _desc = _transceiver.ToString();
 1823        _initialized = true;
 1824        setState(StateNotValidated);
 1825
 1826        return true;
 1827    }
 1828
 1829    private bool validate(int operation)
 1830    {
 1831        if (!_endpoint.datagram()) // Datagram connections are always implicitly validated.
 1832        {
 1833            if (_connector is null) // The server side has the active role for connection validation.
 1834            {
 1835                if (_writeStream.size() == 0)
 1836                {
 1837                    _writeStream.writeBlob(Protocol.magic);
 1838                    ProtocolVersion.ice_write(_writeStream, Protocol.currentProtocol);
 1839                    EncodingVersion.ice_write(_writeStream, Protocol.currentProtocolEncoding);
 1840                    _writeStream.writeByte(Protocol.validateConnectionMsg);
 1841                    _writeStream.writeByte(0); // Compression status (always zero for validate connection).
 1842                    _writeStream.writeInt(Protocol.headerSize); // Message size.
 1843                    TraceUtil.traceSend(_writeStream, _instance, this, _logger, _traceLevels);
 1844                    _writeStream.prepareWrite();
 1845                }
 1846
 1847                if (_observer is not null)
 1848                {
 1849                    observerStartWrite(_writeStream.getBuffer());
 1850                }
 1851
 1852                if (_writeStream.pos() != _writeStream.size())
 1853                {
 1854                    int op = write(_writeStream.getBuffer());
 1855                    if (op != 0)
 1856                    {
 1857                        _threadPool.update(this, operation, op);
 1858                        return false;
 1859                    }
 1860                }
 1861
 1862                if (_observer is not null)
 1863                {
 1864                    observerFinishWrite(_writeStream.getBuffer());
 1865                }
 1866            }
 1867            else // The client side has the passive role for connection validation.
 1868            {
 1869                if (_readStream.size() == 0)
 1870                {
 1871                    _readStream.resize(Protocol.headerSize);
 1872                    _readStream.pos(0);
 1873                }
 1874
 1875                if (_observer is not null)
 1876                {
 1877                    observerStartRead(_readStream.getBuffer());
 1878                }
 1879
 1880                if (_readStream.pos() != _readStream.size())
 1881                {
 1882                    int op = read(_readStream.getBuffer());
 1883                    if (op != 0)
 1884                    {
 1885                        _threadPool.update(this, operation, op);
 1886                        return false;
 1887                    }
 1888                }
 1889
 1890                if (_observer is not null)
 1891                {
 1892                    observerFinishRead(_readStream.getBuffer());
 1893                }
 1894
 1895                _validated = true;
 1896
 1897                Debug.Assert(_readStream.pos() == Protocol.headerSize);
 1898                _readStream.pos(0);
 1899                byte[] m = _readStream.readBlob(4);
 1900                if (m[0] != Protocol.magic[0] || m[1] != Protocol.magic[1] ||
 1901                   m[2] != Protocol.magic[2] || m[3] != Protocol.magic[3])
 1902                {
 1903                    throw new ProtocolException(
 1904                        $"Bad magic in message header: {m[0]:X2} {m[1]:X2} {m[2]:X2} {m[3]:X2}");
 1905                }
 1906
 1907                var pv = new ProtocolVersion(_readStream);
 1908                if (pv != Protocol.currentProtocol)
 1909                {
 1910                    throw new MarshalException(
 1911                        $"Invalid protocol version in message header: {pv.major}.{pv.minor}");
 1912                }
 1913                var ev = new EncodingVersion(_readStream);
 1914                if (ev != Protocol.currentProtocolEncoding)
 1915                {
 1916                    throw new MarshalException(
 1917                        $"Invalid protocol encoding version in message header: {ev.major}.{ev.minor}");
 1918                }
 1919
 1920                byte messageType = _readStream.readByte();
 1921                if (messageType != Protocol.validateConnectionMsg)
 1922                {
 1923                    throw new ProtocolException(
 1924                        $"Received message of type {messageType} over a connection that is not yet validated.");
 1925                }
 1926                _readStream.readByte(); // Ignore compression status for validate connection.
 1927                int size = _readStream.readInt();
 1928                if (size != Protocol.headerSize)
 1929                {
 1930                    throw new MarshalException($"Received ValidateConnection message with unexpected size {size}.");
 1931                }
 1932                TraceUtil.traceRecv(_readStream, this, _logger, _traceLevels);
 1933
 1934                // Client connection starts sending heartbeats once it's received the ValidateConnection message.
 1935                _idleTimeoutTransceiver?.scheduleHeartbeat();
 1936            }
 1937        }
 1938
 1939        _writeStream.resize(0);
 1940        _writeStream.pos(0);
 1941
 1942        _readStream.resize(Protocol.headerSize);
 1943        _readStream.pos(0);
 1944        _readHeader = true;
 1945
 1946        if (_instance.traceLevels().network >= 1)
 1947        {
 1948            var s = new StringBuilder();
 1949            if (_endpoint.datagram())
 1950            {
 1951                s.Append("starting to ");
 1952                s.Append(_connector is not null ? "send" : "receive");
 1953                s.Append(' ');
 1954                s.Append(_endpoint.protocol());
 1955                s.Append(" messages\n");
 1956                s.Append(_transceiver.toDetailedString());
 1957            }
 1958            else
 1959            {
 1960                s.Append(_connector is not null ? "established" : "accepted");
 1961                s.Append(' ');
 1962                s.Append(_endpoint.protocol());
 1963                s.Append(" connection\n");
 1964                s.Append(ToString());
 1965            }
 1966            _instance.initializationData().logger.trace(_instance.traceLevels().networkCat, s.ToString());
 1967        }
 1968
 1969        return true;
 1970    }
 1971
 1972    /// <summary>
 1973    /// Sends the next queued messages. This method is called by message() once the message which is being sent
 1974    /// (_sendStreams.First) is fully sent. Before sending the next message, this message is removed from _sendStreams.
 1975    /// If any, its sent callback is also queued in given callback queue.
 1976    /// </summary>
 1977    /// <param name="callbacks">The sent callbacks to call for the messages that were sent.</param>
 1978    /// <returns>The socket operation to register with the thread pool's selector to send the remainder of the pending
 1979    /// message being sent (_sendStreams.First).</returns>
 1980    private int sendNextMessage(out Queue<OutgoingMessage> callbacks)
 1981    {
 1982        callbacks = null;
 1983
 1984        if (_sendStreams.Count == 0)
 1985        {
 1986            // This can occur if no message was being written and the socket write operation was registered with the
 1987            // thread pool (a transceiver read method can request writing data).
 1988            return SocketOperation.None;
 1989        }
 1990        else if (_state == StateClosingPending && _writeStream.pos() == 0)
 1991        {
 1992            // Message wasn't sent, empty the _writeStream, we're not going to send more data because the connection
 1993            // is being closed.
 1994            OutgoingMessage message = _sendStreams.First.Value;
 1995            _writeStream.swap(message.stream);
 1996            return SocketOperation.None;
 1997        }
 1998
 1999        // Assert that the message was fully written.
 2000        Debug.Assert(!_writeStream.isEmpty() && _writeStream.pos() == _writeStream.size());
 2001
 2002        try
 2003        {
 2004            while (true)
 2005            {
 2006                //
 2007                // The message that was being sent is sent. We can swap back the write stream buffer to the
 2008                // outgoing message (required for retry) and queue its sent callback (if any).
 2009                //
 2010                OutgoingMessage message = _sendStreams.First.Value;
 2011                _writeStream.swap(message.stream);
 2012                if (message.sent())
 2013                {
 2014                    callbacks ??= new Queue<OutgoingMessage>();
 2015                    callbacks.Enqueue(message);
 2016                }
 2017                _sendStreams.RemoveFirst();
 2018
 2019                //
 2020                // If there's nothing left to send, we're done.
 2021                //
 2022                if (_sendStreams.Count == 0)
 2023                {
 2024                    break;
 2025                }
 2026
 2027                //
 2028                // If we are in the closed state or if the close is pending, don't continue sending. This can occur if
 2029                // parseMessage (called before sendNextMessage by message()) closes the connection.
 2030                //
 2031                if (_state >= StateClosingPending)
 2032                {
 2033                    return SocketOperation.None;
 2034                }
 2035
 2036                //
 2037                // Otherwise, prepare the next message.
 2038                //
 2039                message = _sendStreams.First.Value;
 2040                Debug.Assert(!message.prepared);
 2041                OutputStream stream = message.stream;
 2042
 2043                message.stream = doCompress(message.stream, message.compress);
 2044                message.stream.prepareWrite();
 2045                message.prepared = true;
 2046
 2047                TraceUtil.traceSend(stream, _instance, this, _logger, _traceLevels);
 2048
 2049                //
 2050                // Send the message.
 2051                //
 2052                _writeStream.swap(message.stream);
 2053                if (_observer is not null)
 2054                {
 2055                    observerStartWrite(_writeStream.getBuffer());
 2056                }
 2057                if (_writeStream.pos() != _writeStream.size())
 2058                {
 2059                    int op = write(_writeStream.getBuffer());
 2060                    if (op != 0)
 2061                    {
 2062                        return op;
 2063                    }
 2064                }
 2065                if (_observer is not null)
 2066                {
 2067                    observerFinishWrite(_writeStream.getBuffer());
 2068                }
 2069
 2070                // If the message was sent right away, loop to send the next queued message.
 2071            }
 2072
 2073            // Once the CloseConnection message is sent, we transition to the StateClosingPending state.
 2074            if (_state == StateClosing && _shutdownInitiated)
 2075            {
 2076                setState(StateClosingPending);
 2077                int op = _transceiver.closing(true, _exception);
 2078                if (op != 0)
 2079                {
 2080                    return op;
 2081                }
 2082            }
 2083        }
 2084        catch (LocalException ex)
 2085        {
 2086            setState(StateClosed, ex);
 2087        }
 2088        return SocketOperation.None;
 2089    }
 2090
 2091    /// <summary>
 2092    /// Sends or queues the given message.
 2093    /// </summary>
 2094    /// <param name="message">The message to send.</param>
 2095    /// <returns>The send status.</returns>
 2096    private int sendMessage(OutgoingMessage message)
 2097    {
 2098        Debug.Assert(_state >= StateActive);
 2099        Debug.Assert(_state < StateClosed);
 2100
 2101        // Some messages are queued for sending. Just adds the message to the send queue and tell the caller that
 2102        // the message was queued.
 2103        if (_sendStreams.Count > 0)
 2104        {
 2105            _sendStreams.AddLast(message);
 2106            return OutgoingAsyncBase.AsyncStatusQueued;
 2107        }
 2108
 2109        // Prepare the message for sending.
 2110        Debug.Assert(!message.prepared);
 2111
 2112        OutputStream stream = message.stream;
 2113
 2114        message.stream = doCompress(stream, message.compress);
 2115        message.stream.prepareWrite();
 2116        message.prepared = true;
 2117
 2118        TraceUtil.traceSend(stream, _instance, this, _logger, _traceLevels);
 2119
 2120        // Send the message without blocking.
 2121        if (_observer is not null)
 2122        {
 2123            observerStartWrite(message.stream.getBuffer());
 2124        }
 2125        int op = write(message.stream.getBuffer());
 2126        if (op == 0)
 2127        {
 2128            // The message was sent so we're done.
 2129
 2130            if (_observer is not null)
 2131            {
 2132                observerFinishWrite(message.stream.getBuffer());
 2133            }
 2134
 2135            int status = OutgoingAsyncBase.AsyncStatusSent;
 2136            if (message.sent())
 2137            {
 2138                // If there's a sent callback, indicate the caller that it should invoke the sent callback.
 2139                status |= OutgoingAsyncBase.AsyncStatusInvokeSentCallback;
 2140            }
 2141
 2142            return status;
 2143        }
 2144
 2145        // The message couldn't be sent right away so we add it to the send stream queue (which is empty) and swap its
 2146        // stream with `_writeStream`. The socket operation returned by the transceiver write is registered with the
 2147        // thread pool. At this point the message() method will take care of sending the whole message (held by
 2148        // _writeStream) when the transceiver is ready to write more of the message buffer.
 2149
 2150        _writeStream.swap(message.stream);
 2151        _sendStreams.AddLast(message);
 2152        _threadPool.register(this, op);
 2153        return OutgoingAsyncBase.AsyncStatusQueued;
 2154    }
 2155
 2156    private OutputStream doCompress(OutputStream decompressed, bool compress)
 2157    {
 2158        if (BZip2.isLoaded(_logger) && compress && decompressed.size() >= 100)
 2159        {
 2160            //
 2161            // Do compression.
 2162            //
 2163            Ice.Internal.Buffer cbuf = BZip2.compress(
 2164                decompressed.getBuffer(),
 2165                Protocol.headerSize,
 2166                _compressionLevel);
 2167            if (cbuf is not null)
 2168            {
 2169                var cstream = new OutputStream(new Internal.Buffer(cbuf, true), decompressed.getEncoding());
 2170
 2171                //
 2172                // Set compression status.
 2173                //
 2174                cstream.pos(9);
 2175                cstream.writeByte(2);
 2176
 2177                //
 2178                // Write the size of the compressed stream into the header.
 2179                //
 2180                cstream.pos(10);
 2181                cstream.writeInt(cstream.size());
 2182
 2183                //
 2184                // Write the compression status and size of the compressed stream into the header of the
 2185                // decompressed stream -- we need this to trace requests correctly.
 2186                //
 2187                decompressed.pos(9);
 2188                decompressed.writeByte(2);
 2189                decompressed.writeInt(cstream.size());
 2190
 2191                return cstream;
 2192            }
 2193        }
 2194
 2195        // Write the compression status. If BZip2 is loaded and compress is set to true, we write 1, to request a
 2196        // compressed reply. Otherwise, we write 0 either BZip2 is not loaded or we are sending an uncompressed reply.
 2197        decompressed.pos(9);
 2198        decompressed.writeByte((byte)((BZip2.isLoaded(_logger) && compress) ? 1 : 0));
 2199
 2200        //
 2201        // Not compressed, fill in the message size.
 2202        //
 2203        decompressed.pos(10);
 2204        decompressed.writeInt(decompressed.size());
 2205
 2206        return decompressed;
 2207    }
 2208
 2209    private struct MessageInfo
 2210    {
 2211        public InputStream stream;
 2212        public int requestCount;
 2213        public int requestId;
 2214        public byte compress;
 2215        public ObjectAdapter adapter;
 2216        public OutgoingAsyncBase outAsync;
 2217        public int upcallCount;
 2218    }
 2219
 2220    private int parseMessage(ref MessageInfo info)
 2221    {
 2222        Debug.Assert(_state > StateNotValidated && _state < StateClosed);
 2223
 2224        info.stream = new InputStream(_instance, Protocol.currentProtocolEncoding);
 2225        _readStream.swap(info.stream);
 2226        _readStream.resize(Protocol.headerSize);
 2227        _readStream.pos(0);
 2228        _readHeader = true;
 2229
 2230        Debug.Assert(info.stream.pos() == info.stream.size());
 2231
 2232        try
 2233        {
 2234            //
 2235            // The magic and version fields have already been checked.
 2236            //
 2237            info.stream.pos(8);
 2238            byte messageType = info.stream.readByte();
 2239            info.compress = info.stream.readByte();
 2240            if (info.compress == 2)
 2241            {
 2242                if (BZip2.isLoaded(_logger))
 2243                {
 2244                    Ice.Internal.Buffer ubuf = BZip2.decompress(
 2245                        info.stream.getBuffer(),
 2246                        Protocol.headerSize,
 2247                        _messageSizeMax);
 2248                    info.stream = new InputStream(info.stream.instance(), info.stream.getEncoding(), ubuf, true);
 2249                }
 2250                else
 2251                {
 2252                    throw new FeatureNotSupportedException(
 2253                        "Cannot decompress compressed message: BZip2 library is not loaded.");
 2254                }
 2255            }
 2256            info.stream.pos(Protocol.headerSize);
 2257
 2258            switch (messageType)
 2259            {
 2260                case Protocol.closeConnectionMsg:
 2261                {
 2262                    TraceUtil.traceRecv(info.stream, this, _logger, _traceLevels);
 2263                    if (_endpoint.datagram())
 2264                    {
 2265                        if (_warn)
 2266                        {
 2267                            _logger.warning("ignoring close connection message for datagram connection:\n" + _desc);
 2268                        }
 2269                    }
 2270                    else
 2271                    {
 2272                        setState(StateClosingPending, new CloseConnectionException());
 2273
 2274                        //
 2275                        // Notify the transceiver of the graceful connection closure.
 2276                        //
 2277                        int op = _transceiver.closing(false, _exception);
 2278                        if (op != 0)
 2279                        {
 2280                            scheduleCloseTimer();
 2281                            return op;
 2282                        }
 2283                        setState(StateClosed);
 2284                    }
 2285                    break;
 2286                }
 2287
 2288                case Protocol.requestMsg:
 2289                {
 2290                    if (_state >= StateClosing)
 2291                    {
 2292                        TraceUtil.trace(
 2293                            "received request during closing\n(ignored by server, client will retry)",
 2294                            info.stream,
 2295                            this,
 2296                            _logger,
 2297                            _traceLevels);
 2298                    }
 2299                    else
 2300                    {
 2301                        TraceUtil.traceRecv(info.stream, this, _logger, _traceLevels);
 2302                        info.requestId = info.stream.readInt();
 2303                        info.requestCount = 1;
 2304                        info.adapter = _adapter;
 2305                        ++info.upcallCount;
 2306
 2307                        cancelInactivityTimer();
 2308                        ++_dispatchCount;
 2309                    }
 2310                    break;
 2311                }
 2312
 2313                case Protocol.requestBatchMsg:
 2314                {
 2315                    if (_state >= StateClosing)
 2316                    {
 2317                        TraceUtil.trace(
 2318                            "received batch request during closing\n(ignored by server, client will retry)",
 2319                            info.stream,
 2320                            this,
 2321                            _logger,
 2322                            _traceLevels);
 2323                    }
 2324                    else
 2325                    {
 2326                        TraceUtil.traceRecv(info.stream, this, _logger, _traceLevels);
 2327                        int requestCount = info.stream.readInt();
 2328                        if (requestCount < 0)
 2329                        {
 2330                            throw new MarshalException($"Received batch request with {requestCount} batches.");
 2331                        }
 2332
 2333                        // A batched request occupies at least 12 bytes on the wire (a 2-byte identity, a 1-byte
 2334                        // facet path, a 1-byte operation name, a 1-byte operation mode, a 1-byte context, and a
 2335                        // 6-byte parameters encapsulation). Reject a count larger than the remaining message
 2336                        // data could possibly hold. The message size is already capped at Ice.MessageSizeMax
 2337                        // (<= int.MaxValue bytes), so this also keeps requestCount well within range when it is
 2338                        // accumulated into the dispatch counters below.
 2339                        const int minBatchRequestSize = 12;
 2340                        if (requestCount > (info.stream.size() - info.stream.pos()) / minBatchRequestSize)
 2341                        {
 2342                            throw new MarshalException(
 2343                                $"Received batch request with {requestCount} batches, more than the message can contain.
 2344                        }
 2345
 2346                        info.requestCount = requestCount;
 2347                        info.adapter = _adapter;
 2348                        info.upcallCount += info.requestCount;
 2349
 2350                        cancelInactivityTimer();
 2351                        _dispatchCount += info.requestCount;
 2352                    }
 2353                    break;
 2354                }
 2355
 2356                case Protocol.replyMsg:
 2357                {
 2358                    TraceUtil.traceRecv(info.stream, this, _logger, _traceLevels);
 2359                    info.requestId = info.stream.readInt();
 2360                    if (_asyncRequests.TryGetValue(info.requestId, out info.outAsync))
 2361                    {
 2362                        _asyncRequests.Remove(info.requestId);
 2363
 2364                        info.outAsync.getIs().swap(info.stream);
 2365
 2366                        //
 2367                        // If we just received the reply for a request which isn't acknowledge as
 2368                        // sent yet, we queue the reply instead of processing it right away. It
 2369                        // will be processed once the write callback is invoked for the message.
 2370                        //
 2371                        OutgoingMessage message = _sendStreams.Count > 0 ? _sendStreams.First.Value : null;
 2372                        if (message is not null && message.outAsync == info.outAsync)
 2373                        {
 2374                            message.receivedReply = true;
 2375                        }
 2376                        else if (info.outAsync.response())
 2377                        {
 2378                            ++info.upcallCount;
 2379                        }
 2380                        else
 2381                        {
 2382                            info.outAsync = null;
 2383                        }
 2384                        if (_closeRequested && _state < StateClosing && _asyncRequests.Count == 0)
 2385                        {
 2386                            doApplicationClose();
 2387                        }
 2388                    }
 2389                    break;
 2390                }
 2391
 2392                case Protocol.validateConnectionMsg:
 2393                {
 2394                    TraceUtil.traceRecv(info.stream, this, _logger, _traceLevels);
 2395                    break;
 2396                }
 2397
 2398                default:
 2399                {
 2400                    TraceUtil.trace(
 2401                        "received unknown message\n(invalid, closing connection)",
 2402                        info.stream,
 2403                        this,
 2404                        _logger,
 2405                        _traceLevels);
 2406
 2407                    throw new ProtocolException($"Received Ice protocol message with unknown type: {messageType}");
 2408                }
 2409            }
 2410        }
 2411        catch (LocalException ex)
 2412        {
 2413            if (_endpoint.datagram())
 2414            {
 2415                if (_warn)
 2416                {
 2417                    _logger.warning("datagram connection exception:\n" + ex.ToString() + "\n" + _desc);
 2418                }
 2419            }
 2420            else
 2421            {
 2422                setState(StateClosed, ex);
 2423            }
 2424        }
 2425
 2426        if (_state == StateHolding)
 2427        {
 2428            // Don't continue reading if the connection is in the holding state.
 2429            return SocketOperation.None;
 2430        }
 2431        else if (_maxDispatches > 0 && _dispatchCount >= _maxDispatches)
 2432        {
 2433            // Don't continue reading if the _maxDispatches limit is reached or exceeded.
 2434            _idleTimeoutTransceiver?.disableIdleCheck();
 2435            return SocketOperation.None;
 2436        }
 2437        else
 2438        {
 2439            // Continue reading.
 2440            return SocketOperation.Read;
 2441        }
 2442    }
 2443
 2444    private void dispatchAll(
 2445        InputStream stream,
 2446        int requestCount,
 2447        int requestId,
 2448        byte compress,
 2449        ObjectAdapter adapter)
 2450    {
 2451        // Note: In contrast to other private or protected methods, this method must be called *without* the mutex
 2452        // locked.
 2453
 2454        Object dispatcher = adapter?.dispatchPipeline;
 2455
 2456        try
 2457        {
 2458            while (requestCount > 0)
 2459            {
 2460                // adapter can be null here, however the adapter set in current can't be null, and we never pass
 2461                // a null current.adapter to the application code. Once this file enables nullable, adapter should be
 2462                // adapter! below.
 2463                var request = new IncomingRequest(requestId, this, adapter, stream);
 2464
 2465                if (dispatcher is not null)
 2466                {
 2467                    // We don't and can't await the dispatchAsync: with batch requests, we want all the dispatches to
 2468                    // execute in the current Ice thread pool thread. If we awaited the dispatchAsync, we could
 2469                    // switch to a .NET thread pool thread.
 2470                    _ = dispatchAsync(request);
 2471                }
 2472                else
 2473                {
 2474                    // Received request on a connection without an object adapter.
 2475                    sendResponse(
 2476                        request.current.createOutgoingResponse(new ObjectNotExistException()),
 2477                        isTwoWay: !_endpoint.datagram() && requestId != 0,
 2478                        compress: 0);
 2479                }
 2480                --requestCount;
 2481            }
 2482
 2483            stream.clear();
 2484        }
 2485        catch (LocalException ex) // TODO: catch all exceptions
 2486        {
 2487            // Typically, the IncomingRequest constructor throws an exception, and we can't continue.
 2488            dispatchException(ex, requestCount);
 2489        }
 2490
 2491        async Task dispatchAsync(IncomingRequest request)
 2492        {
 2493            try
 2494            {
 2495                OutgoingResponse response;
 2496
 2497                try
 2498                {
 2499                    response = await dispatcher.dispatchAsync(request).ConfigureAwait(false);
 2500                }
 2501                catch (System.Exception ex)
 2502                {
 2503                    response = request.current.createOutgoingResponse(ex);
 2504                }
 2505
 2506                sendResponse(response, isTwoWay: !_endpoint.datagram() && requestId != 0, compress);
 2507            }
 2508            catch (LocalException ex) // TODO: catch all exceptions to avoid UnobservedTaskException
 2509            {
 2510                dispatchException(ex, requestCount: 1);
 2511            }
 2512        }
 2513    }
 2514
 2515    private void sendResponse(OutgoingResponse response, bool isTwoWay, byte compress)
 2516    {
 2517        bool finished = false;
 2518        try
 2519        {
 2520            lock (_mutex)
 2521            {
 2522                Debug.Assert(_state > StateNotValidated);
 2523
 2524                try
 2525                {
 2526                    if (--_upcallCount == 0)
 2527                    {
 2528                        if (_state == StateFinished)
 2529                        {
 2530                            finished = true;
 2531                            _observer?.detach();
 2532                        }
 2533                        Monitor.PulseAll(_mutex);
 2534                    }
 2535
 2536                    if (_state >= StateClosed)
 2537                    {
 2538                        Debug.Assert(_exception is not null);
 2539                        throw _exception;
 2540                    }
 2541
 2542                    if (isTwoWay)
 2543                    {
 2544                        sendMessage(new OutgoingMessage(response.outputStream, compress > 0));
 2545                    }
 2546
 2547                    if (_state == StateActive && _maxDispatches > 0 && _dispatchCount == _maxDispatches)
 2548                    {
 2549                        // Resume reading if the connection is active and the dispatch count is about to be less than
 2550                        // _maxDispatches.
 2551                        _threadPool.update(this, SocketOperation.None, SocketOperation.Read);
 2552                        _idleTimeoutTransceiver?.enableIdleCheck();
 2553                    }
 2554
 2555                    --_dispatchCount;
 2556
 2557                    if (_state == StateClosing && _upcallCount == 0)
 2558                    {
 2559                        initiateShutdown();
 2560                    }
 2561                }
 2562                catch (LocalException ex)
 2563                {
 2564                    setState(StateClosed, ex);
 2565                }
 2566            }
 2567        }
 2568        finally
 2569        {
 2570            if (finished && _removeFromFactory is not null)
 2571            {
 2572                _removeFromFactory(this);
 2573            }
 2574        }
 2575    }
 2576
 2577    private void dispatchException(LocalException ex, int requestCount)
 2578    {
 2579        bool finished = false;
 2580
 2581        // Fatal exception while dispatching a request. Since sendResponse isn't called in case of a fatal exception
 2582        // we decrement _upcallCount here.
 2583        lock (_mutex)
 2584        {
 2585            setState(StateClosed, ex);
 2586
 2587            if (requestCount > 0)
 2588            {
 2589                Debug.Assert(_upcallCount >= requestCount);
 2590                _upcallCount -= requestCount;
 2591                if (_upcallCount == 0)
 2592                {
 2593                    if (_state == StateFinished)
 2594                    {
 2595                        finished = true;
 2596                        _observer?.detach();
 2597                    }
 2598                    Monitor.PulseAll(_mutex);
 2599                }
 2600            }
 2601        }
 2602
 2603        if (finished && _removeFromFactory is not null)
 2604        {
 2605            _removeFromFactory(this);
 2606        }
 2607    }
 2608
 2609    private void inactivityCheck(System.Threading.Timer inactivityTimer)
 2610    {
 2611        lock (_mutex)
 2612        {
 2613            // If the timers are different, it means this inactivityTimer is no longer current.
 2614            if (inactivityTimer == _inactivityTimer)
 2615            {
 2616                _inactivityTimer = null;
 2617                inactivityTimer.Dispose(); // non-blocking
 2618
 2619                if (_state == StateActive)
 2620                {
 2621                    setState(
 2622                        StateClosing,
 2623                        new ConnectionClosedException(
 2624                            "Connection closed because it remained inactive for longer than the inactivity timeout.",
 2625                            closedByApplication: false));
 2626                }
 2627            }
 2628            // Else this timer was already canceled and disposed. Nothing to do.
 2629        }
 2630    }
 2631
 2632    private void connectTimedOut(System.Threading.Timer connectTimer)
 2633    {
 2634        lock (_mutex)
 2635        {
 2636            if (_state < StateActive)
 2637            {
 2638                setState(StateClosed, new ConnectTimeoutException());
 2639            }
 2640        }
 2641        // else ignore since we're already connected.
 2642        connectTimer.Dispose();
 2643    }
 2644
 2645    private void closeTimedOut(System.Threading.Timer closeTimer)
 2646    {
 2647        lock (_mutex)
 2648        {
 2649            if (_state < StateClosed)
 2650            {
 2651                // We don't use setState(state, exception) because we want to overwrite the exception set by a
 2652                // graceful closure.
 2653                _exception = new CloseTimeoutException();
 2654                setState(StateClosed);
 2655            }
 2656        }
 2657        // else ignore since we're already closed.
 2658        closeTimer.Dispose();
 2659    }
 2660
 2661    private ConnectionInfo initConnectionInfo()
 2662    {
 2663        // Called with _mutex locked.
 2664
 2665        if (_state > StateNotInitialized && _info is not null) // Update the connection info until it's initialized
 2666        {
 2667            return _info;
 2668        }
 2669
 2670        _info =
 2671            _transceiver.getInfo(incoming: _connector is null, _adapter?.getName() ?? "", _endpoint.connectionId());
 2672        return _info;
 2673    }
 2674
 2675    private void warning(string msg, System.Exception ex) => _logger.warning($"{msg}:\n{ex}\n{_desc}");
 2676
 2677    private void observerStartRead(Ice.Internal.Buffer buf)
 2678    {
 2679        if (_readStreamPos >= 0)
 2680        {
 2681            Debug.Assert(!buf.empty());
 2682            _observer.receivedBytes(buf.b.position() - _readStreamPos);
 2683        }
 2684        _readStreamPos = buf.empty() ? -1 : buf.b.position();
 2685    }
 2686
 2687    private void observerFinishRead(Ice.Internal.Buffer buf)
 2688    {
 2689        if (_readStreamPos == -1)
 2690        {
 2691            return;
 2692        }
 2693        Debug.Assert(buf.b.position() >= _readStreamPos);
 2694        _observer.receivedBytes(buf.b.position() - _readStreamPos);
 2695        _readStreamPos = -1;
 2696    }
 2697
 2698    private void observerStartWrite(Ice.Internal.Buffer buf)
 2699    {
 2700        if (_writeStreamPos >= 0)
 2701        {
 2702            Debug.Assert(!buf.empty());
 2703            _observer.sentBytes(buf.b.position() - _writeStreamPos);
 2704        }
 2705        _writeStreamPos = buf.empty() ? -1 : buf.b.position();
 2706    }
 2707
 2708    private void observerFinishWrite(Ice.Internal.Buffer buf)
 2709    {
 2710        if (_writeStreamPos == -1)
 2711        {
 2712            return;
 2713        }
 2714        if (buf.b.position() > _writeStreamPos)
 2715        {
 2716            _observer.sentBytes(buf.b.position() - _writeStreamPos);
 2717        }
 2718        _writeStreamPos = -1;
 2719    }
 2720
 2721    private int read(Ice.Internal.Buffer buf)
 2722    {
 2723        int start = buf.b.position();
 2724        int op = _transceiver.read(buf, ref _hasMoreData);
 2725        if (_instance.traceLevels().network >= 3 && buf.b.position() != start)
 2726        {
 2727            var s = new StringBuilder("received ");
 2728            if (_endpoint.datagram())
 2729            {
 2730                s.Append(buf.b.limit());
 2731            }
 2732            else
 2733            {
 2734                s.Append(buf.b.position() - start);
 2735                s.Append(" of ");
 2736                s.Append(buf.b.limit() - start);
 2737            }
 2738            s.Append(" bytes via ");
 2739            s.Append(_endpoint.protocol());
 2740            s.Append('\n');
 2741            s.Append(ToString());
 2742            _instance.initializationData().logger.trace(_instance.traceLevels().networkCat, s.ToString());
 2743        }
 2744        return op;
 2745    }
 2746
 2747    private int write(Ice.Internal.Buffer buf)
 2748    {
 2749        int start = buf.b.position();
 2750        int op = _transceiver.write(buf);
 2751        if (_instance.traceLevels().network >= 3 && buf.b.position() != start)
 2752        {
 2753            var s = new StringBuilder("sent ");
 2754            s.Append(buf.b.position() - start);
 2755            if (!_endpoint.datagram())
 2756            {
 2757                s.Append(" of ");
 2758                s.Append(buf.b.limit() - start);
 2759            }
 2760            s.Append(" bytes via ");
 2761            s.Append(_endpoint.protocol());
 2762            s.Append('\n');
 2763            s.Append(ToString());
 2764            _instance.initializationData().logger.trace(_instance.traceLevels().networkCat, s.ToString());
 2765        }
 2766        return op;
 2767    }
 2768
 2769    private void scheduleInactivityTimer()
 2770    {
 2771        // Called with the ConnectionI mutex locked.
 2772        Debug.Assert(_inactivityTimer is null);
 2773        Debug.Assert(_inactivityTimeout > TimeSpan.Zero);
 2774
 2775        _inactivityTimer = new System.Threading.Timer(
 2776            inactivityTimer => inactivityCheck((System.Threading.Timer)inactivityTimer));
 2777        _inactivityTimer.Change(_inactivityTimeout, Timeout.InfiniteTimeSpan);
 2778    }
 2779
 2780    private void cancelInactivityTimer()
 2781    {
 2782        // Called with the ConnectionI mutex locked.
 2783        if (_inactivityTimer is not null)
 2784        {
 2785            _inactivityTimer.Change(Timeout.InfiniteTimeSpan, Timeout.InfiniteTimeSpan);
 2786            _inactivityTimer.Dispose();
 2787            _inactivityTimer = null;
 2788        }
 2789    }
 2790
 2791    private void scheduleCloseTimer()
 2792    {
 2793        if (_closeTimeout > TimeSpan.Zero)
 2794        {
 2795#pragma warning disable CA2000 // closeTimer is disposed by closeTimedOut.
 2796            var closeTimer = new System.Threading.Timer(
 2797                timerObj => closeTimedOut((System.Threading.Timer)timerObj));
 2798            // schedule timer to run once; closeTimedOut disposes the timer too.
 2799            closeTimer.Change(_closeTimeout, Timeout.InfiniteTimeSpan);
 2800#pragma warning restore CA2000
 2801        }
 2802    }
 2803
 2804    private void doApplicationClose()
 2805    {
 2806        // Called with the ConnectionI mutex locked.
 2807        Debug.Assert(_state < StateClosing);
 2808        setState(
 2809            StateClosing,
 2810            new ConnectionClosedException(
 2811                "The connection was closed gracefully by the application.",
 2812                closedByApplication: true));
 2813    }
 2814
 2815    private class OutgoingMessage
 2816    {
 12817        internal OutgoingMessage(OutputStream stream, bool compress)
 2818        {
 12819            this.stream = stream;
 12820            this.compress = compress;
 12821        }
 2822
 12823        internal OutgoingMessage(OutgoingAsyncBase outAsync, OutputStream stream, bool compress, int requestId)
 2824        {
 12825            this.outAsync = outAsync;
 12826            this.stream = stream;
 12827            this.compress = compress;
 12828            this.requestId = requestId;
 12829        }
 2830
 2831        internal void canceled()
 2832        {
 2833            Debug.Assert(outAsync is not null); // Only requests can timeout.
 12834            outAsync = null;
 12835        }
 2836
 2837        internal bool sent()
 2838        {
 12839            stream = null;
 12840            if (outAsync is not null)
 2841            {
 12842                invokeSent = outAsync.sent();
 12843                return invokeSent || receivedReply;
 2844            }
 12845            return false;
 2846        }
 2847
 2848        internal void completed(LocalException ex)
 2849        {
 12850            if (outAsync is not null)
 2851            {
 12852                if (outAsync.exception(ex))
 2853                {
 12854                    outAsync.invokeException();
 2855                }
 2856            }
 12857            stream = null;
 12858        }
 2859
 2860        internal OutputStream stream;
 2861        internal OutgoingAsyncBase outAsync;
 2862        internal bool compress;
 2863        internal int requestId;
 2864        internal bool prepared;
 2865        internal bool isSent;
 2866        internal bool invokeSent;
 2867        internal bool receivedReply;
 2868    }
 2869
 2870    private static readonly ConnectionState[] connectionStateMap = [
 2871        ConnectionState.ConnectionStateValidating,   // StateNotInitialized
 2872        ConnectionState.ConnectionStateValidating,   // StateNotValidated
 2873        ConnectionState.ConnectionStateActive,       // StateActive
 2874        ConnectionState.ConnectionStateHolding,      // StateHolding
 2875        ConnectionState.ConnectionStateClosing,      // StateClosing
 2876        ConnectionState.ConnectionStateClosing,      // StateClosingPending
 2877        ConnectionState.ConnectionStateClosed,       // StateClosed
 2878        ConnectionState.ConnectionStateClosed,       // StateFinished
 2879    ];
 2880
 2881    private readonly Instance _instance;
 2882    private readonly Transceiver _transceiver;
 2883    private readonly IdleTimeoutTransceiverDecorator _idleTimeoutTransceiver; // can be null
 2884
 2885    private string _desc;
 2886    private readonly string _type;
 2887    private readonly Connector _connector;
 2888    private readonly EndpointI _endpoint;
 2889
 2890    private ObjectAdapter _adapter;
 2891
 2892    private readonly Logger _logger;
 2893    private readonly TraceLevels _traceLevels;
 2894    private readonly Ice.Internal.ThreadPool _threadPool;
 2895    private readonly bool _threadHopRequired;
 2896
 2897    private readonly TimeSpan _connectTimeout;
 2898    private readonly TimeSpan _closeTimeout;
 2899    private TimeSpan _inactivityTimeout; // protected by _mutex
 2900
 2901    private System.Threading.Timer _inactivityTimer; // can be null
 2902
 2903    private StartCallback _startCallback;
 2904
 2905    // This action must be called outside the ConnectionI lock to avoid lock acquisition deadlocks.
 2906    private readonly Action<ConnectionI> _removeFromFactory;
 2907
 2908    private readonly bool _warn;
 2909    private readonly bool _warnUdp;
 2910
 2911    private readonly int _compressionLevel;
 2912
 2913    private int _nextRequestId;
 2914
 2915    private readonly Dictionary<int, OutgoingAsyncBase> _asyncRequests = new Dictionary<int, OutgoingAsyncBase>();
 2916
 2917    private LocalException _exception;
 2918
 2919    private readonly int _messageSizeMax;
 2920    private readonly BatchRequestQueue _batchRequestQueue;
 2921
 2922    private readonly LinkedList<OutgoingMessage> _sendStreams = new LinkedList<OutgoingMessage>();
 2923
 2924    // Contains the message which is being received. If the connection is waiting to receive a message (_readHeader ==
 2925    // true), its size is Protocol.headerSize. Otherwise, its size is the message size specified in the received message
 2926    // header.
 2927    private readonly InputStream _readStream;
 2928
 2929    // When _readHeader is true, the next bytes we'll read are the header of a new message. When false, we're reading
 2930    // next the remainder of a message that was already partially received.
 2931    private bool _readHeader;
 2932
 2933    // Contains the message which is being sent. The write stream buffer is empty if no message is being sent.
 2934    private readonly OutputStream _writeStream;
 2935
 2936    private ConnectionObserver _observer;
 2937    private int _readStreamPos;
 2938    private int _writeStreamPos;
 2939
 2940    // The upcall count keeps track of the number of dispatches, AMI (response) continuations, sent callbacks and
 2941    // connection establishment callbacks that have been started (or are about to be started) by a thread of the thread
 2942    // pool associated with this connection, and have not completed yet. All these operations except the connection
 2943    // establishment callbacks execute application code or code generated from Slice definitions.
 2944    private int _upcallCount;
 2945
 2946    // The number of outstanding dispatches. Maintained only while state is StateActive or StateHolding.
 2947    // _dispatchCount can be greater than a non-0 _maxDispatches when a receive a batch with multiples requests.
 2948    private int _dispatchCount;
 2949
 2950    // When we dispatch _maxDispatches concurrent requests, we stop reading the connection to back-pressure the peer.
 2951    // _maxDispatches <= 0 means no limit.
 2952    private readonly int _maxDispatches;
 2953
 2954    private int _state; // The current state.
 2955    private bool _shutdownInitiated;
 2956    private bool _initialized;
 2957    private bool _validated;
 2958
 2959    // When true, the application called close and Connection must close the connection when it receives the reply
 2960    // for the last outstanding invocation.
 2961    private bool _closeRequested;
 2962
 2963    private ConnectionInfo _info;
 2964
 2965    private CloseCallback _closeCallback;
 2966
 2967    // We need to run the continuation asynchronously since it can be completed by an Ice thread pool thread.
 2968    private readonly TaskCompletionSource _closed = new(TaskCreationOptions.RunContinuationsAsynchronously);
 2969    private readonly object _mutex = new();
 2970}