< Summary

Information
Class: Ice.ConnectionI
Assembly: Ice
File(s): /_/csharp/src/Ice/ConnectionI.cs
Tag: 125_37167941578
Line coverage
87%
Covered lines: 1058
Uncovered lines: 158
Coverable lines: 1216
Total lines: 2970
Line coverage: 87%
Branch coverage
81%
Covered branches: 599
Total branches: 732
Branch coverage: 81.8%
Method coverage
94%
Covered methods: 74
Fully covered methods: 46
Total methods: 78
Method coverage: 94.8%
Full method coverage: 58.9%

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
start(...)87.5%8894.44%
startAndWait()30%281043.75%
activate()50%2283.33%
hold()50%2283.33%
destroy(...)50%44100%
abort()100%11100%
closeAsync()100%44100%
isActiveOrHolding()100%22100%
throwException()100%22100%
waitUntilHolding()100%44100%
waitUntilFinished()100%44100%
updateObserver()66.67%6692.86%
sendAsyncRequest(...)90%101093.33%
getBatchRequestQueue()100%11100%
flushBatchRequests(...)100%1175%
flushBatchRequestsAsync(...)100%11100%
disableInactivityCheck()100%11100%
setCloseCallback(...)25%8438.46%
asyncRequestCanceled(...)52.94%473477.42%
endpoint()100%11100%
connector()100%11100%
setAdapter(...)62.5%9880%
getAdapter()100%11100%
getEndpoint()100%11100%
createProxy(...)100%11100%
setAdapterFromAdapter(...)50%4485.71%
startAsync(...)50%5466.67%
doIO()95%202092.86%
finishAsync(...)100%2424100%
message(...)84.15%968287.18%
upcall(...)96.67%303091.89%
finished(...)100%88100%
finish()91.67%626091.43%
ToString()100%11100%
type()100%11100%
getInfo()100%22100%
setBufferSize(...)50%2263.64%
exception(...)100%11100%
getThreadPool()100%210%
.ctor(...)81.25%161696.49%
idleCheck(...)100%44100%
sendHeartbeat()65.38%282685.19%
isHeartbeat()100%210%
toConnectionState(...)100%11100%
setState(...)75%202092.31%
setState(...)87.14%897084.38%
initiateShutdown()100%88100%
initialize(...)100%22100%
validate(...)76%625083.33%
sendNextMessage(...)92.86%312883.72%
sendMessage(...)100%1010100%
doCompress(...)100%88100%
parseMessage(...)73.91%714677.23%
dispatchAll(...)87.5%8881.25%
dispatchAsync()100%2272.73%
sendResponse(...)96.15%2626100%
dispatchException(...)0%156120%
inactivityCheck(...)100%44100%
connectTimedOut(...)100%22100%
closeTimedOut(...)100%22100%
initConnectionInfo()100%88100%
warning(...)100%210%
observerStartRead(...)75%44100%
observerFinishRead(...)50%2280%
observerStartWrite(...)100%44100%
observerFinishWrite(...)100%44100%
read(...)83.33%6693.33%
write(...)100%66100%
scheduleInactivityTimer()100%11100%
cancelInactivityTimer()100%22100%
scheduleCloseTimer()100%22100%
doApplicationClose()100%11100%
.ctor(...)100%11100%
.ctor(...)100%11100%
canceled()100%11100%
sent()100%44100%
completed(...)100%44100%
.cctor()100%11100%

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        {
 125            lock (_mutex)
 26            {
 27                //
 28                // The connection might already be closed if the communicator was destroyed.
 29                //
 130                if (_state >= StateClosed)
 31                {
 32                    Debug.Assert(_exception is not null);
 033                    throw _exception;
 34                }
 35
 136                if (!initialize(SocketOperation.None) || !validate(SocketOperation.None))
 37                {
 138                    if (_connectTimeout > TimeSpan.Zero)
 39                    {
 40#pragma warning disable CA2000 // connectTimer is disposed by connectTimedOut.
 141                        var connectTimer = new System.Threading.Timer(
 142                            timerObj => connectTimedOut((System.Threading.Timer)timerObj));
 43                        // schedule timer to run once; connectTimedOut disposes the timer too.
 144                        connectTimer.Change(_connectTimeout, Timeout.InfiniteTimeSpan);
 45#pragma warning restore CA2000
 46                    }
 47
 148                    _startCallback = callback;
 149                    return;
 50                }
 51
 52                // The connection starts in the holding state. It will be activated by the connection factory.
 153                setState(StateHolding);
 154            }
 155        }
 156        catch (LocalException ex)
 57        {
 158            exception(ex);
 159            callback.connectionStartFailed(this, _exception);
 160            return;
 61        }
 62
 163        callback.connectionStartCompleted(this);
 164    }
 65
 66    internal void startAndWait()
 67    {
 68        try
 69        {
 170            lock (_mutex)
 71            {
 72                //
 73                // The connection might already be closed if the communicator was destroyed.
 74                //
 175                if (_state >= StateClosed)
 76                {
 77                    Debug.Assert(_exception is not null);
 078                    throw _exception;
 79                }
 80
 181                if (!initialize(SocketOperation.None) || !validate(SocketOperation.None))
 82                {
 83                    //
 84                    // Wait for the connection to be validated.
 85                    //
 086                    while (_state <= StateNotValidated)
 87                    {
 088                        Monitor.Wait(_mutex);
 89                    }
 90
 091                    if (_state >= StateClosing)
 92                    {
 93                        Debug.Assert(_exception is not null);
 094                        throw _exception;
 95                    }
 96                }
 97
 98                //
 99                // We start out in holding state.
 100                //
 1101                setState(StateHolding);
 1102            }
 1103        }
 0104        catch (LocalException ex)
 105        {
 0106            exception(ex);
 0107            waitUntilFinished();
 0108            return;
 109        }
 1110    }
 111
 112    internal void activate()
 113    {
 1114        lock (_mutex)
 115        {
 1116            if (_state <= StateNotValidated)
 117            {
 0118                return;
 119            }
 120
 1121            setState(StateActive);
 1122        }
 1123    }
 124
 125    internal void hold()
 126    {
 1127        lock (_mutex)
 128        {
 1129            if (_state <= StateNotValidated)
 130            {
 0131                return;
 132            }
 133
 1134            setState(StateHolding);
 1135        }
 1136    }
 137
 138    // DestructionReason.
 139    public const int ObjectAdapterDeactivated = 0;
 140    public const int CommunicatorDestroyed = 1;
 141
 142    internal void destroy(int reason)
 143    {
 1144        lock (_mutex)
 145        {
 146            switch (reason)
 147            {
 148                case ObjectAdapterDeactivated:
 149                {
 1150                    setState(StateClosing, new ObjectAdapterDeactivatedException(_adapter?.getName() ?? ""));
 1151                    break;
 152                }
 153
 154                case CommunicatorDestroyed:
 155                {
 1156                    setState(StateClosing, new CommunicatorDestroyedException());
 1157                    break;
 158                }
 159            }
 160        }
 1161    }
 162
 163    public void abort()
 164    {
 1165        lock (_mutex)
 166        {
 1167            setState(
 1168                StateClosed,
 1169                new ConnectionAbortedException(
 1170                    "The connection was aborted by the application.",
 1171                    closedByApplication: true));
 1172        }
 1173    }
 174
 175    public Task closeAsync()
 176    {
 1177        lock (_mutex)
 178        {
 1179            if (_state < StateClosing)
 180            {
 1181                if (_asyncRequests.Count == 0)
 182                {
 1183                    doApplicationClose();
 184                }
 185                else
 186                {
 1187                    _closeRequested = true;
 1188                    scheduleCloseTimer(); // we don't wait forever for outstanding invocations to complete
 189                }
 190            }
 191            // else nothing to do, already closing or closed.
 1192        }
 193
 1194        return _closed.Task;
 195    }
 196
 197    internal bool isActiveOrHolding()
 198    {
 1199        lock (_mutex)
 200        {
 1201            return _state > StateNotValidated && _state < StateClosing;
 202        }
 1203    }
 204
 205    public void throwException()
 206    {
 1207        lock (_mutex)
 208        {
 1209            if (_exception is not null)
 210            {
 211                Debug.Assert(_state >= StateClosing);
 1212                throw _exception;
 213            }
 1214        }
 1215    }
 216
 217    internal void waitUntilHolding()
 218    {
 1219        lock (_mutex)
 220        {
 1221            while (_state < StateHolding || _upcallCount > 0)
 222            {
 1223                Monitor.Wait(_mutex);
 224            }
 1225        }
 1226    }
 227
 228    internal void waitUntilFinished()
 229    {
 1230        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            //
 1238            while (_state < StateFinished || _upcallCount > 0)
 239            {
 1240                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            //
 1248            _adapter = null;
 1249        }
 1250    }
 251
 252    internal void updateObserver()
 253    {
 1254        lock (_mutex)
 255        {
 1256            if (_state < StateNotValidated || _state > StateClosed)
 257            {
 0258                return;
 259            }
 260
 261            Debug.Assert(_instance.initializationData().observer is not null);
 1262            _observer = _instance.initializationData().observer.getConnectionObserver(
 1263                initConnectionInfo(),
 1264                _endpoint,
 1265                toConnectionState(_state),
 1266                _observer);
 1267            if (_observer is not null)
 268            {
 1269                _observer.attach();
 270            }
 271            else
 272            {
 1273                _writeStreamPos = -1;
 1274                _readStreamPos = -1;
 275            }
 1276        }
 1277    }
 278
 279    internal int sendAsyncRequest(
 280        OutgoingAsyncBase og,
 281        bool compress,
 282        bool response,
 283        int batchRequestCount)
 284    {
 1285        OutputStream os = og.getOs();
 286
 1287        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            //
 1294            if (_exception is not null)
 295            {
 1296                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            //
 1306            _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            //
 1312            og.cancelable(this);
 1313            int requestId = 0;
 1314            if (response)
 315            {
 316                //
 317                // Create a new unique request ID.
 318                //
 1319                requestId = _nextRequestId++;
 1320                if (requestId <= 0)
 321                {
 0322                    _nextRequestId = 1;
 0323                    requestId = _nextRequestId++;
 324                }
 325
 326                //
 327                // Fill in the request ID.
 328                //
 1329                os.pos(Protocol.headerSize);
 1330                os.writeInt(requestId);
 331            }
 1332            else if (batchRequestCount > 0)
 333            {
 1334                os.pos(Protocol.headerSize);
 1335                os.writeInt(batchRequestCount);
 336            }
 337
 1338            og.attachRemoteObserver(initConnectionInfo(), _endpoint, requestId);
 339
 340            // We're just about to send a request, so we are not inactive anymore.
 1341            cancelInactivityTimer();
 342
 1343            int status = OutgoingAsyncBase.AsyncStatusQueued;
 344            try
 345            {
 1346                var message = new OutgoingMessage(og, os, compress, requestId);
 1347                status = sendMessage(message);
 1348            }
 1349            catch (LocalException ex)
 350            {
 1351                setState(StateClosed, ex);
 352                Debug.Assert(_exception is not null);
 1353                throw _exception;
 354            }
 355
 1356            if (response)
 357            {
 358                //
 359                // Add to the async requests map.
 360                //
 1361                _asyncRequests[requestId] = og;
 362            }
 1363            return status;
 364        }
 1365    }
 366
 1367    internal BatchRequestQueue getBatchRequestQueue() => _batchRequestQueue;
 368
 369    public void flushBatchRequests(CompressBatch compress)
 370    {
 371        try
 372        {
 1373            var completed = new FlushBatchTaskCompletionCallback();
 1374            var outgoing = new ConnectionFlushBatchAsync(this, _instance, completed);
 1375            outgoing.invoke(_flushBatchRequests_name, compress, true);
 1376            completed.Task.Wait();
 1377        }
 0378        catch (AggregateException ex)
 379        {
 0380            throw ex.InnerException;
 381        }
 1382    }
 383
 384    public Task flushBatchRequestsAsync(
 385        CompressBatch compress,
 386        IProgress<bool> progress = null,
 387        CancellationToken cancel = default)
 388    {
 1389        var completed = new FlushBatchTaskCompletionCallback(progress, cancel);
 1390        var outgoing = new ConnectionFlushBatchAsync(this, _instance, completed);
 1391        outgoing.invoke(_flushBatchRequests_name, compress, false);
 1392        return completed.Task;
 393    }
 394
 395    private const string _flushBatchRequests_name = "flushBatchRequests";
 396
 397    public void disableInactivityCheck()
 398    {
 1399        lock (_mutex)
 400        {
 1401            cancelInactivityTimer();
 1402            _inactivityTimeout = TimeSpan.Zero;
 1403        }
 1404    }
 405
 406    public void setCloseCallback(CloseCallback callback)
 407    {
 1408        lock (_mutex)
 409        {
 1410            if (_state >= StateClosed)
 411            {
 0412                if (callback is not null)
 413                {
 0414                    _threadPool.execute(
 0415                        () =>
 0416                        {
 0417                            try
 0418                            {
 0419                                callback(this);
 0420                            }
 0421                            catch (System.Exception ex)
 0422                            {
 0423                                _logger.error("connection callback exception:\n" + ex + '\n' + _desc);
 0424                            }
 0425                        },
 0426                        this);
 427                }
 428            }
 429            else
 430            {
 1431                _closeCallback = callback;
 432            }
 1433        }
 1434    }
 435
 436    public void asyncRequestCanceled(OutgoingAsyncBase outAsync, LocalException ex)
 437    {
 438        //
 439        // NOTE: This isn't called from a thread pool thread.
 440        //
 441
 1442        lock (_mutex)
 443        {
 1444            if (_state >= StateClosed)
 445            {
 0446                return; // The request has already been or will be shortly notified of the failure.
 447            }
 448
 1449            OutgoingMessage o = _sendStreams.FirstOrDefault(m => m.outAsync == outAsync);
 1450            if (o is not null)
 451            {
 1452                if (o.requestId > 0)
 453                {
 1454                    _asyncRequests.Remove(o.requestId);
 455                }
 456
 1457                if (ex is ConnectionAbortedException)
 458                {
 0459                    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                    //
 1467                    if (o == _sendStreams.First.Value)
 468                    {
 0469                        o.canceled();
 470                    }
 471                    else
 472                    {
 1473                        o.canceled();
 1474                        _sendStreams.Remove(o);
 475                    }
 1476                    if (outAsync.exception(ex))
 477                    {
 1478                        outAsync.invokeExceptionAsync();
 479                    }
 480                }
 481
 1482                if (_closeRequested && _state < StateClosing && _asyncRequests.Count == 0)
 483                {
 0484                    doApplicationClose();
 485                }
 1486                return;
 487            }
 488
 1489            if (outAsync is OutgoingAsync)
 490            {
 1491                foreach (KeyValuePair<int, OutgoingAsyncBase> kvp in _asyncRequests)
 492                {
 1493                    if (kvp.Value == outAsync)
 494                    {
 1495                        if (ex is ConnectionAbortedException)
 496                        {
 0497                            setState(StateClosed, ex);
 498                        }
 499                        else
 500                        {
 1501                            _asyncRequests.Remove(kvp.Key);
 1502                            if (outAsync.exception(ex))
 503                            {
 1504                                outAsync.invokeExceptionAsync();
 505                            }
 506                        }
 507
 1508                        if (_closeRequested && _state < StateClosing && _asyncRequests.Count == 0)
 509                        {
 0510                            doApplicationClose();
 511                        }
 1512                        return;
 513                    }
 514                }
 515            }
 0516        }
 1517    }
 518
 1519    internal EndpointI endpoint() => _endpoint; // No mutex protection necessary, _endpoint is immutable.
 520
 1521    internal Connector connector() => _connector; // No mutex protection necessary, _endpoint is immutable.
 522
 523    public void setAdapter(ObjectAdapter adapter)
 524    {
 1525        if (_connector is null) // server connection
 526        {
 0527            throw new InvalidOperationException("setAdapter can only be called on a client connection");
 528        }
 529
 1530        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.
 1534            adapter.setAdapterOnConnection(this);
 535        }
 536        else
 537        {
 1538            lock (_mutex)
 539            {
 1540                if (_state <= StateNotValidated || _state >= StateClosing)
 541                {
 0542                    return;
 543                }
 1544                _adapter = null;
 1545            }
 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        //
 1552    }
 553
 554    public ObjectAdapter getAdapter()
 555    {
 1556        lock (_mutex)
 557        {
 1558            return _adapter;
 559        }
 1560    }
 561
 1562    public Endpoint getEndpoint() => _endpoint; // No mutex protection necessary, _endpoint is immutable.
 563
 564    public ObjectPrx createProxy(Identity id)
 565    {
 1566        ObjectAdapter.checkIdentity(id);
 1567        return new ObjectPrxHelper(_instance.referenceFactory().create(id, this));
 568    }
 569
 570    public void setAdapterFromAdapter(ObjectAdapter adapter)
 571    {
 1572        lock (_mutex)
 573        {
 1574            if (_state <= StateNotValidated || _state >= StateClosing)
 575            {
 0576                return;
 577            }
 578            Debug.Assert(adapter is not null); // Called by ObjectAdapter::setAdapterOnConnection
 1579            _adapter = adapter;
 580
 581            // Clear cached connection info (if any) as it's no longer accurate.
 1582            _info = null;
 1583        }
 1584    }
 585
 586    //
 587    // Operations from EventHandler
 588    //
 589    public override bool startAsync(int operation, Ice.Internal.AsyncCallback completedCallback)
 590    {
 1591        if (_state >= StateClosed)
 592        {
 0593            return false;
 594        }
 595
 1596        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.
 0600            Task.Run(doIO);
 601        }
 602        else
 603        {
 1604            doIO();
 605        }
 606
 1607        return true;
 608
 609        void doIO()
 610        {
 1611            lock (_mutex)
 612            {
 1613                if (_state >= StateClosed)
 614                {
 0615                    completedCallback(this);
 0616                    return;
 617                }
 618
 619                try
 620                {
 1621                    if ((operation & SocketOperation.Write) != 0)
 622                    {
 1623                        if (_observer != null)
 624                        {
 1625                            observerStartWrite(_writeStream.getBuffer());
 626                        }
 627
 1628                        bool completedSynchronously =
 1629                            _transceiver.startWrite(
 1630                                _writeStream.getBuffer(),
 1631                                completedCallback,
 1632                                this,
 1633                                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.
 1637                        if (messageWritten && _sendStreams.Count > 0)
 638                        {
 639                            // See finish() code.
 1640                            _sendStreams.First.Value.isSent = true;
 641                        }
 642
 1643                        if (completedSynchronously)
 644                        {
 645                            // If the write completed synchronously, we need to call the completedCallback.
 1646                            completedCallback(this);
 647                        }
 648                    }
 1649                    else if ((operation & SocketOperation.Read) != 0)
 650                    {
 1651                        if (_observer != null && !_readHeader)
 652                        {
 1653                            observerStartRead(_readStream.getBuffer());
 654                        }
 655
 1656                        if (_transceiver.startRead(_readStream.getBuffer(), completedCallback, this))
 657                        {
 1658                            completedCallback(this);
 659                        }
 660                    }
 1661                }
 1662                catch (LocalException ex)
 663                {
 1664                    setState(StateClosed, ex);
 1665                    completedCallback(this);
 1666                }
 667            }
 1668        }
 669    }
 670
 671    public override bool finishAsync(int operation)
 672    {
 1673        if (_state >= StateClosed)
 674        {
 1675            return false;
 676        }
 677
 678        try
 679        {
 1680            if ((operation & SocketOperation.Write) != 0)
 681            {
 1682                Ice.Internal.Buffer buf = _writeStream.getBuffer();
 1683                int start = buf.b.position();
 1684                _transceiver.finishWrite(buf);
 1685                if (_instance.traceLevels().network >= 3 && buf.b.position() != start)
 686                {
 1687                    var s = new StringBuilder("sent ");
 1688                    s.Append(buf.b.position() - start);
 1689                    if (!_endpoint.datagram())
 690                    {
 1691                        s.Append(" of ");
 1692                        s.Append(buf.b.limit() - start);
 693                    }
 1694                    s.Append(" bytes via ");
 1695                    s.Append(_endpoint.protocol());
 1696                    s.Append('\n');
 1697                    s.Append(ToString());
 1698                    _instance.initializationData().logger.trace(_instance.traceLevels().networkCat, s.ToString());
 699                }
 700
 1701                if (_observer is not null)
 702                {
 1703                    observerFinishWrite(_writeStream.getBuffer());
 704                }
 705            }
 1706            else if ((operation & SocketOperation.Read) != 0)
 707            {
 1708                Ice.Internal.Buffer buf = _readStream.getBuffer();
 1709                int start = buf.b.position();
 1710                _transceiver.finishRead(buf);
 1711                if (_instance.traceLevels().network >= 3 && buf.b.position() != start)
 712                {
 1713                    var s = new StringBuilder("received ");
 1714                    if (_endpoint.datagram())
 715                    {
 1716                        s.Append(buf.b.limit());
 717                    }
 718                    else
 719                    {
 1720                        s.Append(buf.b.position() - start);
 1721                        s.Append(" of ");
 1722                        s.Append(buf.b.limit() - start);
 723                    }
 1724                    s.Append(" bytes via ");
 1725                    s.Append(_endpoint.protocol());
 1726                    s.Append('\n');
 1727                    s.Append(ToString());
 1728                    _instance.initializationData().logger.trace(_instance.traceLevels().networkCat, s.ToString());
 729                }
 730
 1731                if (_observer is not null && !_readHeader)
 732                {
 1733                    observerFinishRead(_readStream.getBuffer());
 734                }
 735            }
 1736        }
 1737        catch (LocalException ex)
 738        {
 1739            setState(StateClosed, ex);
 1740        }
 1741        return _state < StateClosed;
 742    }
 743
 744    public override void message(ThreadPoolCurrent current)
 745    {
 1746        StartCallback startCB = null;
 1747        Queue<OutgoingMessage> sentCBs = null;
 1748        var info = new MessageInfo();
 1749        int upcallCount = 0;
 750
 1751        using var msg = new ThreadPoolMessage(current, _mutex);
 1752        lock (_mutex)
 753        {
 754            try
 755            {
 1756                if (!msg.startIOScope())
 757                {
 1758                    return;
 759                }
 760
 1761                if (_state >= StateClosed)
 762                {
 0763                    return;
 764                }
 765
 766                try
 767                {
 1768                    int writeOp = SocketOperation.None;
 1769                    int readOp = SocketOperation.None;
 770
 771                    // If writes are ready, write the data from the connection's write buffer (_writeStream)
 1772                    if ((current.operation & SocketOperation.Write) != 0)
 773                    {
 1774                        if (_observer is not null)
 775                        {
 1776                            observerStartWrite(_writeStream.getBuffer());
 777                        }
 1778                        writeOp = write(_writeStream.getBuffer());
 1779                        if (_observer is not null && (writeOp & SocketOperation.Write) == 0)
 780                        {
 1781                            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
 1790                    if ((current.operation & SocketOperation.Read) != 0)
 791                    {
 792                        while (true)
 793                        {
 1794                            Ice.Internal.Buffer buf = _readStream.getBuffer();
 795
 1796                            if (_observer is not null && !_readHeader)
 797                            {
 1798                                observerStartRead(buf);
 799                            }
 800
 1801                            readOp = read(buf);
 1802                            if ((readOp & SocketOperation.Read) != 0)
 803                            {
 804                                // Can't continue without blocking, exit out of the loop.
 805                                break;
 806                            }
 1807                            if (_observer is not null && !_readHeader)
 808                            {
 809                                Debug.Assert(!buf.b.hasRemaining());
 1810                                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.
 1815                            if (_readHeader)
 816                            {
 817                                // The next read will read the remainder of the message.
 1818                                _readHeader = false;
 819
 1820                                _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                                //
 1829                                _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).
 1833                                int pos = _readStream.pos();
 1834                                if (pos < Protocol.headerSize)
 835                                {
 836                                    //
 837                                    // This situation is possible for small UDP packets.
 838                                    //
 0839                                    throw new MarshalException("Received Ice message with too few bytes in header.");
 840                                }
 841
 842                                // Decode the header.
 1843                                _readStream.pos(0);
 1844                                byte[] m = new byte[4];
 1845                                m[0] = _readStream.readByte();
 1846                                m[1] = _readStream.readByte();
 1847                                m[2] = _readStream.readByte();
 1848                                m[3] = _readStream.readByte();
 1849                                if (m[0] != Protocol.magic[0] || m[1] != Protocol.magic[1] ||
 1850                                m[2] != Protocol.magic[2] || m[3] != Protocol.magic[3])
 851                                {
 0852                                    throw new ProtocolException(
 0853                                        $"Bad magic in message header: {m[0]:X2} {m[1]:X2} {m[2]:X2} {m[3]:X2}");
 854                                }
 855
 1856                                var pv = new ProtocolVersion(_readStream);
 1857                                if (pv != Protocol.currentProtocol)
 858                                {
 0859                                    throw new MarshalException(
 0860                                        $"Invalid protocol version in message header: {pv.major}.{pv.minor}");
 861                                }
 1862                                var ev = new EncodingVersion(_readStream);
 1863                                if (ev != Protocol.currentProtocolEncoding)
 864                                {
 0865                                    throw new MarshalException(
 0866                                        $"Invalid protocol encoding version in message header: {ev.major}.{ev.minor}");
 867                                }
 868
 1869                                _readStream.readByte(); // messageType
 1870                                _readStream.readByte(); // compress
 1871                                int size = _readStream.readInt();
 1872                                if (size < Protocol.headerSize)
 873                                {
 0874                                    throw new MarshalException($"Received Ice message with unexpected size {size}.");
 875                                }
 876
 877                                // Resize the read buffer to the message size.
 1878                                if (size > _messageSizeMax)
 879                                {
 1880                                    Ex.throwMemoryLimitException(size, _messageSizeMax);
 881                                }
 1882                                if (size > _readStream.size())
 883                                {
 1884                                    _readStream.resize(size);
 885                                }
 1886                                _readStream.pos(pos);
 887                            }
 888
 1889                            if (buf.b.hasRemaining())
 890                            {
 1891                                if (_endpoint.datagram())
 892                                {
 1893                                    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.
 1904                    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.
 1909                    int readyOp = current.operation & ~newOp;
 910
 1911                    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.
 1915                        if (newOp != 0)
 916                        {
 1917                            _threadPool.update(this, current.operation, newOp);
 1918                            return;
 919                        }
 920
 921                        // Initialize the connection if it's not initialized yet.
 1922                        if (_state == StateNotInitialized && !initialize(current.operation))
 923                        {
 1924                            return;
 925                        }
 926
 927                        // Validate the connection if it's not validated yet.
 1928                        if (_state <= StateNotValidated && !validate(current.operation))
 929                        {
 1930                            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.
 1935                        _threadPool.unregister(this, current.operation);
 936
 937                        //
 938                        // We start out in holding state.
 939                        //
 1940                        setState(StateHolding);
 1941                        if (_startCallback is not null)
 942                        {
 1943                            startCB = _startCallback;
 1944                            _startCallback = null;
 1945                            ++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                        //
 1956                        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.
 1960                            newOp |= parseMessage(ref info);
 1961                            upcallCount += info.upcallCount;
 962                        }
 963
 1964                        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
 1969                            newOp |= sendNextMessage(out sentCBs);
 1970                            if (sentCBs is not null)
 971                            {
 1972                                ++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.
 1978                        if (_state < StateClosed)
 979                        {
 1980                            _threadPool.update(this, current.operation, newOp);
 981                        }
 982                    }
 983
 1984                    if (upcallCount == 0)
 985                    {
 1986                        return; // Nothing to execute, we're done!
 987                    }
 988
 1989                    _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.
 1993                    msg.ioCompleted();
 1994                }
 1995                catch (DatagramLimitException) // Expected.
 996                {
 1997                    if (_warnUdp)
 998                    {
 0999                        _logger.warning($"maximum datagram size of {_readStream.pos()} exceeded");
 1000                    }
 11001                    _readStream.resize(Protocol.headerSize);
 11002                    _readStream.pos(0);
 11003                    _readHeader = true;
 11004                    return;
 1005                }
 11006                catch (SocketException ex)
 1007                {
 11008                    setState(StateClosed, ex);
 11009                    return;
 1010                }
 11011                catch (LocalException ex)
 1012                {
 11013                    if (_endpoint.datagram())
 1014                    {
 01015                        if (_warn)
 1016                        {
 01017                            _logger.warning($"datagram connection exception:\n{ex}\n{_desc}");
 1018                        }
 01019                        _readStream.resize(Protocol.headerSize);
 01020                        _readStream.pos(0);
 01021                        _readHeader = true;
 1022                    }
 1023                    else
 1024                    {
 11025                        setState(StateClosed, ex);
 1026                    }
 11027                    return;
 1028                }
 1029            }
 1030            finally
 1031            {
 11032                msg.finishIOScope();
 11033            }
 1034        }
 1035
 11036        _threadPool.executeFromThisThread(() => upcall(startCB, sentCBs, info), this);
 11037    }
 1038
 1039    private void upcall(StartCallback startCB, Queue<OutgoingMessage> sentCBs, MessageInfo info)
 1040    {
 11041        int completedUpcallCount = 0;
 1042
 1043        //
 1044        // Notify the factory that the connection establishment and
 1045        // validation has completed.
 1046        //
 11047        if (startCB is not null)
 1048        {
 11049            startCB.connectionStartCompleted(this);
 11050            ++completedUpcallCount;
 1051        }
 1052
 1053        //
 1054        // Notify AMI calls that the message was sent.
 1055        //
 11056        if (sentCBs is not null)
 1057        {
 11058            foreach (OutgoingMessage m in sentCBs)
 1059            {
 11060                if (m.invokeSent)
 1061                {
 11062                    m.outAsync.invokeSent();
 1063                }
 11064                if (m.receivedReply)
 1065                {
 11066                    var outAsync = (OutgoingAsync)m.outAsync;
 11067                    if (outAsync.response())
 1068                    {
 11069                        outAsync.invokeResponse();
 1070                    }
 1071                }
 1072            }
 11073            ++completedUpcallCount;
 1074        }
 1075
 1076        //
 1077        // Asynchronous replies must be handled outside the thread
 1078        // synchronization, so that nested calls are possible.
 1079        //
 11080        if (info.outAsync is not null)
 1081        {
 11082            info.outAsync.invokeResponse();
 11083            ++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        //
 11091        if (info.requestCount > 0)
 1092        {
 11093            dispatchAll(info.stream, info.requestCount, info.requestId, info.compress, info.adapter);
 1094        }
 1095
 1096        //
 1097        // Decrease the upcall count.
 1098        //
 11099        bool finished = false;
 11100        if (completedUpcallCount > 0)
 1101        {
 11102            lock (_mutex)
 1103            {
 11104                _upcallCount -= completedUpcallCount;
 11105                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.
 11109                    if (_state == StateClosing)
 1110                    {
 1111                        try
 1112                        {
 11113                            initiateShutdown();
 11114                        }
 01115                        catch (Ice.LocalException ex)
 1116                        {
 01117                            setState(StateClosed, ex);
 01118                        }
 1119                    }
 11120                    else if (_state == StateFinished)
 1121                    {
 11122                        finished = true;
 11123                        _observer?.detach();
 1124                    }
 11125                    Monitor.PulseAll(_mutex);
 1126                }
 11127            }
 1128        }
 1129
 11130        if (finished && _removeFromFactory is not null)
 1131        {
 11132            _removeFromFactory(this);
 1133        }
 11134    }
 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).
 11142        lock (_mutex)
 1143        {
 1144            Debug.Assert(_state == StateClosed);
 11145        }
 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        //
 11152        if (_startCallback is null && _sendStreams.Count == 0 && _asyncRequests.Count == 0 && _closeCallback is null)
 1153        {
 11154            finish();
 11155            return;
 1156        }
 1157
 11158        current.ioCompleted();
 11159        _threadPool.executeFromThisThread(finish, this);
 11160    }
 1161
 1162    private void finish()
 1163    {
 11164        if (!_initialized)
 1165        {
 11166            if (_instance.traceLevels().network >= 2)
 1167            {
 11168                var s = new StringBuilder("failed to ");
 11169                s.Append(_connector is not null ? "establish" : "accept");
 11170                s.Append(' ');
 11171                s.Append(_endpoint.protocol());
 11172                s.Append(" connection\n");
 11173                s.Append(ToString());
 11174                s.Append('\n');
 11175                s.Append(_exception);
 11176                _instance.initializationData().logger.trace(_instance.traceLevels().networkCat, s.ToString());
 1177            }
 1178        }
 1179        else
 1180        {
 11181            if (_instance.traceLevels().network >= 1)
 1182            {
 11183                var s = new StringBuilder("closed ");
 11184                s.Append(_endpoint.protocol());
 11185                s.Append(" connection\n");
 11186                s.Append(ToString());
 1187
 1188                // Trace the cause of most connection closures.
 11189                if (!(_exception is CommunicatorDestroyedException || _exception is ObjectAdapterDeactivatedException))
 1190                {
 11191                    s.Append('\n');
 11192                    s.Append(_exception);
 1193                }
 1194
 11195                _instance.initializationData().logger.trace(_instance.traceLevels().networkCat, s.ToString());
 1196            }
 1197        }
 1198
 11199        _startCallback?.connectionStartFailed(this, _exception);
 11200        _startCallback = null;
 1201
 11202        if (_sendStreams.Count > 0)
 1203        {
 11204            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                //
 11210                OutgoingMessage message = _sendStreams.First.Value;
 11211                _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                //
 11218                if (message.isSent || message.receivedReply)
 1219                {
 11220                    if (message.sent() && message.invokeSent)
 1221                    {
 11222                        message.outAsync.invokeSent();
 1223                    }
 11224                    if (message.receivedReply)
 1225                    {
 01226                        var outAsync = (OutgoingAsync)message.outAsync;
 01227                        if (outAsync.response())
 1228                        {
 01229                            outAsync.invokeResponse();
 1230                        }
 1231                    }
 11232                    _sendStreams.RemoveFirst();
 1233                }
 1234            }
 1235
 11236            foreach (OutgoingMessage o in _sendStreams)
 1237            {
 11238                o.completed(_exception);
 11239                if (o.requestId > 0) // Make sure finished isn't called twice.
 1240                {
 11241                    _asyncRequests.Remove(o.requestId);
 1242                }
 1243            }
 11244            _sendStreams.Clear(); // Must be cleared before _requests because of Outgoing* references in OutgoingMessage
 1245        }
 1246
 11247        foreach (OutgoingAsyncBase o in _asyncRequests.Values)
 1248        {
 11249            if (o.exception(_exception))
 1250            {
 11251                o.invokeException();
 1252            }
 1253        }
 11254        _asyncRequests.Clear();
 1255
 1256        //
 1257        // Don't wait to be reaped to reclaim memory allocated by read/write streams.
 1258        //
 11259        _writeStream.clear();
 11260        _writeStream.getBuffer().clear();
 11261        _readStream.clear();
 11262        _readStream.getBuffer().clear();
 1263
 11264        if (_exception is ConnectionClosedException or
 11265            CloseConnectionException or
 11266            CommunicatorDestroyedException or
 11267            ObjectAdapterDeactivatedException)
 1268        {
 1269            // Can execute synchronously. Note that we're not within a lock(this) here.
 11270            _closed.SetResult();
 1271        }
 1272        else
 1273        {
 1274            Debug.Assert(_exception is not null);
 11275            _closed.SetException(_exception);
 1276        }
 1277
 11278        if (_closeCallback is not null)
 1279        {
 1280            try
 1281            {
 11282                _closeCallback(this);
 11283            }
 01284            catch (System.Exception ex)
 1285            {
 01286                _logger.error("connection callback exception:\n" + ex + '\n' + _desc);
 01287            }
 11288            _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        //
 11295        bool finished = false;
 11296        lock (_mutex)
 1297        {
 11298            setState(StateFinished);
 1299
 11300            if (_upcallCount == 0)
 1301            {
 11302                finished = true;
 11303                _observer?.detach();
 1304            }
 11305        }
 1306
 11307        if (finished && _removeFromFactory is not null)
 1308        {
 11309            _removeFromFactory(this);
 1310        }
 11311    }
 1312
 1313    /// <inheritdoc/>
 11314    public override string ToString() => _desc; // No mutex lock, _desc is immutable.
 1315
 1316    /// <inheritdoc/>
 11317    public string type() => _type; // No mutex lock, _type is immutable.
 1318
 1319    /// <inheritdoc/>
 1320    public ConnectionInfo getInfo()
 1321    {
 11322        lock (_mutex)
 1323        {
 11324            if (_state >= StateClosed)
 1325            {
 11326                throw _exception;
 1327            }
 11328            return initConnectionInfo();
 1329        }
 11330    }
 1331
 1332    /// <inheritdoc/>
 1333    public void setBufferSize(int rcvSize, int sndSize)
 1334    {
 11335        lock (_mutex)
 1336        {
 11337            if (_state >= StateClosed)
 1338            {
 01339                throw _exception;
 1340            }
 1341            try
 1342            {
 11343                _transceiver.setBufferSize(rcvSize, sndSize);
 11344            }
 01345            catch (LocalException ex)
 1346            {
 1347                // The failing call may have closed the socket, so close the connection as well.
 01348                setState(StateClosed, ex);
 01349                throw;
 1350            }
 11351            _info = null; // Invalidate the cached connection info
 11352        }
 11353    }
 1354
 1355    public void exception(LocalException ex)
 1356    {
 11357        lock (_mutex)
 1358        {
 11359            setState(StateClosed, ex);
 11360        }
 11361    }
 1362
 01363    public Ice.Internal.ThreadPool getThreadPool() => _threadPool;
 1364
 11365    internal ConnectionI(
 11366        Instance instance,
 11367        Transceiver transceiver,
 11368        Connector connector, // null for incoming connections, non-null for outgoing connections
 11369        EndpointI endpoint,
 11370        ObjectAdapter adapter,
 11371        Action<ConnectionI> removeFromFactory, // can be null
 11372        ConnectionOptions options)
 1373    {
 11374        _instance = instance;
 11375        _desc = transceiver.ToString();
 11376        _type = transceiver.protocol();
 11377        _connector = connector;
 11378        _endpoint = endpoint;
 11379        _adapter = adapter;
 11380        InitializationData initData = instance.initializationData();
 11381        _logger = initData.logger; // Cached for better performance.
 11382        _traceLevels = instance.traceLevels(); // Cached for better performance.
 11383        _connectTimeout = options.connectTimeout;
 11384        _closeTimeout = options.closeTimeout; // not used for datagram connections
 1385        // suppress inactivity timeout for datagram connections
 11386        _inactivityTimeout = endpoint.datagram() ? TimeSpan.Zero : options.inactivityTimeout;
 11387        _maxDispatches = options.maxDispatches;
 11388        _removeFromFactory = removeFromFactory;
 11389        _warn = initData.properties.getIcePropertyAsInt("Ice.Warn.Connections") > 0;
 11390        _warnUdp = initData.properties.getIcePropertyAsInt("Ice.Warn.Datagrams") > 0;
 11391        _nextRequestId = 1;
 11392        _messageSizeMax = connector is null ? adapter.messageSizeMax() : instance.messageSizeMax();
 11393        _batchRequestQueue = new BatchRequestQueue(instance, _endpoint.datagram());
 11394        _readStream = new InputStream(instance, Protocol.currentProtocolEncoding);
 11395        _readHeader = false;
 11396        _readStreamPos = -1;
 11397        _writeStream = new OutputStream(); // temporary stream
 11398        _writeStreamPos = -1;
 11399        _upcallCount = 0;
 11400        _state = StateNotInitialized;
 1401
 11402        _compressionLevel = initData.properties.getIcePropertyAsInt("Ice.Compression.Level");
 11403        if (_compressionLevel < 1)
 1404        {
 01405            _compressionLevel = 1;
 1406        }
 11407        else if (_compressionLevel > 9)
 1408        {
 01409            _compressionLevel = 9;
 1410        }
 1411
 11412        if (options.idleTimeout > TimeSpan.Zero && !endpoint.datagram())
 1413        {
 11414            _idleTimeoutTransceiver = new IdleTimeoutTransceiverDecorator(
 11415                transceiver,
 11416                this,
 11417                options.idleTimeout,
 11418                options.enableIdleCheck);
 11419            transceiver = _idleTimeoutTransceiver;
 1420        }
 11421        _transceiver = transceiver;
 1422
 11423        if (connector is null)
 1424        {
 1425            // adapter is always set for incoming connections
 1426            Debug.Assert(adapter is not null);
 11427            _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.
 11433            _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.
 11440        _threadHopRequired = AssemblyUtil.isWindows && _threadPool.canShrink;
 1441
 1442        // initialize only resets the handler state; it doesn't throw.
 11443        _threadPool.initialize(this);
 11444    }
 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    {
 11452        lock (_mutex)
 1453        {
 11454            if (_state == StateActive && _idleTimeoutTransceiver!.idleCheckEnabled)
 1455            {
 11456                int idleTimeoutInSeconds = (int)idleTimeout.TotalSeconds;
 1457
 11458                setState(
 11459                    StateClosed,
 11460                    new ConnectionAbortedException(
 11461                        $"Connection aborted by the idle check because it did not receive any bytes for {idleTimeoutInSe
 11462                        closedByApplication: false));
 1463            }
 1464            // else nothing to do
 11465        }
 11466    }
 1467
 1468    internal void sendHeartbeat()
 1469    {
 1470        Debug.Assert(!_endpoint.datagram());
 1471
 11472        lock (_mutex)
 1473        {
 11474            if (_state == StateActive || _state == StateHolding || _state == StateClosing)
 1475            {
 1476                // We check if the connection has become inactive.
 11477                if (
 11478                    _inactivityTimer is null &&           // timer not already scheduled
 11479                    _inactivityTimeout > TimeSpan.Zero && // inactivity timeout is enabled
 11480                    _state == StateActive &&              // only schedule the timer if the connection is active
 11481                    _dispatchCount == 0 &&                // no pending dispatch
 11482                    _asyncRequests.Count == 0 &&          // no pending invocation
 11483                    _readHeader &&                        // we're not waiting for the remainder of an incoming message
 11484                    _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.
 11491                    if (_sendStreams.Count == 0 || isHeartbeat(_writeStream))
 1492                    {
 11493                        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.
 11504                if (_sendStreams.Count == 0)
 1505                {
 11506                    var os = new OutputStream(Protocol.currentProtocolEncoding);
 11507                    os.writeBlob(Protocol.magic);
 11508                    ProtocolVersion.ice_write(os, Protocol.currentProtocol);
 11509                    EncodingVersion.ice_write(os, Protocol.currentProtocolEncoding);
 11510                    os.writeByte(Protocol.validateConnectionMsg);
 11511                    os.writeByte(0);
 11512                    os.writeInt(Protocol.headerSize); // Message size.
 1513                    try
 1514                    {
 11515                        _ = sendMessage(new OutgoingMessage(os, compress: false));
 11516                    }
 01517                    catch (LocalException ex)
 1518                    {
 01519                        setState(StateClosed, ex);
 01520                    }
 1521                }
 1522            }
 1523            // else nothing to do
 01524        }
 1525
 1526        static bool isHeartbeat(OutputStream stream) =>
 01527            stream.getBuffer().b.get(8) == Protocol.validateConnectionMsg;
 11528    }
 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
 11539    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
 11549        if (_state == state) // Don't switch twice.
 1550        {
 11551            return;
 1552        }
 1553
 11554        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
 11561            _exception = ex;
 1562
 1563            //
 1564            // We don't warn if we are not validated.
 1565            //
 11566            if (_warn && _validated)
 1567            {
 1568                //
 1569                // Don't warn about certain expected exceptions.
 1570                //
 11571                if (!(_exception is CloseConnectionException ||
 11572                     _exception is ConnectionClosedException ||
 11573                     _exception is CommunicatorDestroyedException ||
 11574                     _exception is ObjectAdapterDeactivatedException ||
 11575                     (_exception is ConnectionLostException && _state >= StateClosing)))
 1576                {
 01577                    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        //
 11587        setState(state);
 11588    }
 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        //
 11596        if (_endpoint.datagram() && state == StateClosing)
 1597        {
 11598            state = StateClosed;
 1599        }
 1600
 1601        //
 1602        // Skip graceful shutdown if we are destroyed before validation.
 1603        //
 11604        if (_state <= StateNotValidated && state == StateClosing)
 1605        {
 11606            state = StateClosed;
 1607        }
 1608
 11609        if (_state == state) // Don't switch twice.
 1610        {
 01611            return;
 1612        }
 1613
 11614        if (state > StateActive)
 1615        {
 1616            // Dispose the inactivity timer, if not null.
 11617            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                {
 11632                    if (_state != StateNotInitialized)
 1633                    {
 1634                        Debug.Assert(_state == StateClosed);
 01635                        return;
 1636                    }
 1637                    break;
 1638                }
 1639
 1640                case StateActive:
 1641                {
 1642                    //
 1643                    // Can only switch to active from holding or not validated.
 1644                    //
 11645                    if (_state != StateHolding && _state != StateNotValidated)
 1646                    {
 01647                        return;
 1648                    }
 1649
 11650                    if (_maxDispatches <= 0 || _dispatchCount < _maxDispatches)
 1651                    {
 11652                        _threadPool.register(this, SocketOperation.Read);
 11653                        _idleTimeoutTransceiver?.enableIdleCheck();
 1654                    }
 1655                    // else don't resume reading since we're at or over the _maxDispatches limit.
 1656
 11657                    break;
 1658                }
 1659
 1660                case StateHolding:
 1661                {
 1662                    //
 1663                    // Can only switch to holding from active or not validated.
 1664                    //
 11665                    if (_state != StateActive && _state != StateNotValidated)
 1666                    {
 01667                        return;
 1668                    }
 1669
 11670                    if (_state == StateActive && (_maxDispatches <= 0 || _dispatchCount < _maxDispatches))
 1671                    {
 11672                        _threadPool.unregister(this, SocketOperation.Read);
 11673                        _idleTimeoutTransceiver?.disableIdleCheck();
 1674                    }
 1675                    // else reads are already disabled because the _maxDispatches limit is reached or exceeded.
 1676
 11677                    break;
 1678                }
 1679
 1680                case StateClosing:
 1681                case StateClosingPending:
 1682                {
 1683                    //
 1684                    // Can't change back from closing pending.
 1685                    //
 11686                    if (_state >= StateClosingPending)
 1687                    {
 11688                        return;
 1689                    }
 1690                    break;
 1691                }
 1692
 1693                case StateClosed:
 1694                {
 11695                    if (_state == StateFinished)
 1696                    {
 11697                        return;
 1698                    }
 1699
 11700                    _batchRequestQueue.destroy(_exception);
 11701                    _threadPool.finish(this);
 11702                    _transceiver.close();
 11703                    break;
 1704                }
 1705
 1706                case StateFinished:
 1707                {
 1708                    Debug.Assert(_state == StateClosed);
 11709                    _transceiver.destroy();
 1710                    break;
 1711                }
 1712            }
 11713        }
 01714        catch (LocalException ex)
 1715        {
 01716            _logger.error("unexpected connection exception:\n" + ex + "\n" + _desc);
 01717        }
 1718
 11719        if (_instance.initializationData().observer is not null)
 1720        {
 11721            ConnectionState oldState = toConnectionState(_state);
 11722            ConnectionState newState = toConnectionState(state);
 11723            if (oldState != newState)
 1724            {
 11725                _observer = _instance.initializationData().observer.getConnectionObserver(
 11726                    initConnectionInfo(),
 11727                    _endpoint,
 11728                    newState,
 11729                    _observer);
 11730                if (_observer is not null)
 1731                {
 11732                    _observer.attach();
 1733                }
 1734                else
 1735                {
 11736                    _writeStreamPos = -1;
 11737                    _readStreamPos = -1;
 1738                }
 1739            }
 11740            if (_observer is not null && state == StateClosed && _exception is not null)
 1741            {
 11742                if (!(_exception is CloseConnectionException ||
 11743                     _exception is ConnectionClosedException ||
 11744                     _exception is CommunicatorDestroyedException ||
 11745                     _exception is ObjectAdapterDeactivatedException ||
 11746                     (_exception is ConnectionLostException && _state >= StateClosing)))
 1747                {
 11748                    _observer.failed(_exception.ice_id());
 1749                }
 1750            }
 1751        }
 11752        _state = state;
 1753
 11754        Monitor.PulseAll(_mutex);
 1755
 11756        if (_state == StateClosing && _upcallCount == 0)
 1757        {
 1758            try
 1759            {
 11760                initiateShutdown();
 11761            }
 01762            catch (LocalException ex)
 1763            {
 01764                setState(StateClosed, ex);
 01765            }
 1766        }
 11767    }
 1768
 1769    private void initiateShutdown()
 1770    {
 1771        Debug.Assert(_state == StateClosing && _upcallCount == 0);
 1772
 11773        if (_shutdownInitiated)
 1774        {
 11775            return;
 1776        }
 11777        _shutdownInitiated = true;
 1778
 11779        if (!_endpoint.datagram())
 1780        {
 1781            //
 1782            // Before we shut down, we send a close connection message.
 1783            //
 11784            var os = new OutputStream(Protocol.currentProtocolEncoding);
 11785            os.writeBlob(Protocol.magic);
 11786            ProtocolVersion.ice_write(os, Protocol.currentProtocol);
 11787            EncodingVersion.ice_write(os, Protocol.currentProtocolEncoding);
 11788            os.writeByte(Protocol.closeConnectionMsg);
 11789            os.writeByte(0); // Compression status: always zero for close connection.
 11790            os.writeInt(Protocol.headerSize); // Message size.
 1791
 11792            scheduleCloseTimer();
 1793
 11794            if ((sendMessage(new OutgoingMessage(os, compress: false)) & OutgoingAsyncBase.AsyncStatusSent) != 0)
 1795            {
 11796                setState(StateClosingPending);
 1797
 1798                //
 1799                // Notify the transceiver of the graceful connection closure.
 1800                //
 11801                int op = _transceiver.closing(true, _exception);
 11802                if (op != 0)
 1803                {
 11804                    _threadPool.register(this, op);
 1805                }
 1806            }
 1807        }
 11808    }
 1809
 1810    private bool initialize(int operation)
 1811    {
 11812        int s = _transceiver.initialize(_readStream.getBuffer(), _writeStream.getBuffer(), ref _hasMoreData);
 11813        if (s != SocketOperation.None)
 1814        {
 11815            _threadPool.update(this, operation, s);
 11816            return false;
 1817        }
 1818
 1819        //
 1820        // Update the connection description once the transceiver is initialized.
 1821        //
 11822        _desc = _transceiver.ToString();
 11823        _initialized = true;
 11824        setState(StateNotValidated);
 1825
 11826        return true;
 1827    }
 1828
 1829    private bool validate(int operation)
 1830    {
 11831        if (!_endpoint.datagram()) // Datagram connections are always implicitly validated.
 1832        {
 11833            if (_connector is null) // The server side has the active role for connection validation.
 1834            {
 11835                if (_writeStream.size() == 0)
 1836                {
 11837                    _writeStream.writeBlob(Protocol.magic);
 11838                    ProtocolVersion.ice_write(_writeStream, Protocol.currentProtocol);
 11839                    EncodingVersion.ice_write(_writeStream, Protocol.currentProtocolEncoding);
 11840                    _writeStream.writeByte(Protocol.validateConnectionMsg);
 11841                    _writeStream.writeByte(0); // Compression status (always zero for validate connection).
 11842                    _writeStream.writeInt(Protocol.headerSize); // Message size.
 11843                    TraceUtil.traceSend(_writeStream, _instance, this, _logger, _traceLevels);
 11844                    _writeStream.prepareWrite();
 1845                }
 1846
 11847                if (_observer is not null)
 1848                {
 01849                    observerStartWrite(_writeStream.getBuffer());
 1850                }
 1851
 11852                if (_writeStream.pos() != _writeStream.size())
 1853                {
 11854                    int op = write(_writeStream.getBuffer());
 11855                    if (op != 0)
 1856                    {
 11857                        _threadPool.update(this, operation, op);
 11858                        return false;
 1859                    }
 1860                }
 1861
 11862                if (_observer is not null)
 1863                {
 01864                    observerFinishWrite(_writeStream.getBuffer());
 1865                }
 1866            }
 1867            else // The client side has the passive role for connection validation.
 1868            {
 11869                if (_readStream.size() == 0)
 1870                {
 11871                    _readStream.resize(Protocol.headerSize);
 11872                    _readStream.pos(0);
 1873                }
 1874
 11875                if (_observer is not null)
 1876                {
 01877                    observerStartRead(_readStream.getBuffer());
 1878                }
 1879
 11880                if (_readStream.pos() != _readStream.size())
 1881                {
 11882                    int op = read(_readStream.getBuffer());
 11883                    if (op != 0)
 1884                    {
 11885                        _threadPool.update(this, operation, op);
 11886                        return false;
 1887                    }
 1888                }
 1889
 11890                if (_observer is not null)
 1891                {
 01892                    observerFinishRead(_readStream.getBuffer());
 1893                }
 1894
 11895                _validated = true;
 1896
 1897                Debug.Assert(_readStream.pos() == Protocol.headerSize);
 11898                _readStream.pos(0);
 11899                byte[] m = _readStream.readBlob(4);
 11900                if (m[0] != Protocol.magic[0] || m[1] != Protocol.magic[1] ||
 11901                   m[2] != Protocol.magic[2] || m[3] != Protocol.magic[3])
 1902                {
 01903                    throw new ProtocolException(
 01904                        $"Bad magic in message header: {m[0]:X2} {m[1]:X2} {m[2]:X2} {m[3]:X2}");
 1905                }
 1906
 11907                var pv = new ProtocolVersion(_readStream);
 11908                if (pv != Protocol.currentProtocol)
 1909                {
 01910                    throw new MarshalException(
 01911                        $"Invalid protocol version in message header: {pv.major}.{pv.minor}");
 1912                }
 11913                var ev = new EncodingVersion(_readStream);
 11914                if (ev != Protocol.currentProtocolEncoding)
 1915                {
 01916                    throw new MarshalException(
 01917                        $"Invalid protocol encoding version in message header: {ev.major}.{ev.minor}");
 1918                }
 1919
 11920                byte messageType = _readStream.readByte();
 11921                if (messageType != Protocol.validateConnectionMsg)
 1922                {
 01923                    throw new ProtocolException(
 01924                        $"Received message of type {messageType} over a connection that is not yet validated.");
 1925                }
 11926                _readStream.readByte(); // Ignore compression status for validate connection.
 11927                int size = _readStream.readInt();
 11928                if (size != Protocol.headerSize)
 1929                {
 01930                    throw new MarshalException($"Received ValidateConnection message with unexpected size {size}.");
 1931                }
 11932                TraceUtil.traceRecv(_readStream, this, _logger, _traceLevels);
 1933
 1934                // Client connection starts sending heartbeats once it's received the ValidateConnection message.
 11935                _idleTimeoutTransceiver?.scheduleHeartbeat();
 1936            }
 1937        }
 1938
 11939        _writeStream.resize(0);
 11940        _writeStream.pos(0);
 1941
 11942        _readStream.resize(Protocol.headerSize);
 11943        _readStream.pos(0);
 11944        _readHeader = true;
 1945
 11946        if (_instance.traceLevels().network >= 1)
 1947        {
 11948            var s = new StringBuilder();
 11949            if (_endpoint.datagram())
 1950            {
 11951                s.Append("starting to ");
 11952                s.Append(_connector is not null ? "send" : "receive");
 11953                s.Append(' ');
 11954                s.Append(_endpoint.protocol());
 11955                s.Append(" messages\n");
 11956                s.Append(_transceiver.toDetailedString());
 1957            }
 1958            else
 1959            {
 11960                s.Append(_connector is not null ? "established" : "accepted");
 11961                s.Append(' ');
 11962                s.Append(_endpoint.protocol());
 11963                s.Append(" connection\n");
 11964                s.Append(ToString());
 1965            }
 11966            _instance.initializationData().logger.trace(_instance.traceLevels().networkCat, s.ToString());
 1967        }
 1968
 11969        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    {
 11982        callbacks = null;
 1983
 11984        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).
 11988            return SocketOperation.None;
 1989        }
 11990        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.
 01994            OutgoingMessage message = _sendStreams.First.Value;
 01995            _writeStream.swap(message.stream);
 01996            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        {
 12004            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                //
 12010                OutgoingMessage message = _sendStreams.First.Value;
 12011                _writeStream.swap(message.stream);
 12012                if (message.sent())
 2013                {
 12014                    callbacks ??= new Queue<OutgoingMessage>();
 12015                    callbacks.Enqueue(message);
 2016                }
 12017                _sendStreams.RemoveFirst();
 2018
 2019                //
 2020                // If there's nothing left to send, we're done.
 2021                //
 12022                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                //
 12031                if (_state >= StateClosingPending)
 2032                {
 02033                    return SocketOperation.None;
 2034                }
 2035
 2036                //
 2037                // Otherwise, prepare the next message.
 2038                //
 12039                message = _sendStreams.First.Value;
 2040                Debug.Assert(!message.prepared);
 12041                OutputStream stream = message.stream;
 2042
 12043                message.stream = doCompress(message.stream, message.compress);
 12044                message.stream.prepareWrite();
 12045                message.prepared = true;
 2046
 12047                TraceUtil.traceSend(stream, _instance, this, _logger, _traceLevels);
 2048
 2049                //
 2050                // Send the message.
 2051                //
 12052                _writeStream.swap(message.stream);
 12053                if (_observer is not null)
 2054                {
 12055                    observerStartWrite(_writeStream.getBuffer());
 2056                }
 12057                if (_writeStream.pos() != _writeStream.size())
 2058                {
 12059                    int op = write(_writeStream.getBuffer());
 12060                    if (op != 0)
 2061                    {
 12062                        return op;
 2063                    }
 2064                }
 12065                if (_observer is not null)
 2066                {
 12067                    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.
 12074            if (_state == StateClosing && _shutdownInitiated)
 2075            {
 12076                setState(StateClosingPending);
 12077                int op = _transceiver.closing(true, _exception);
 12078                if (op != 0)
 2079                {
 12080                    return op;
 2081                }
 2082            }
 12083        }
 02084        catch (LocalException ex)
 2085        {
 02086            setState(StateClosed, ex);
 02087        }
 12088        return SocketOperation.None;
 12089    }
 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.
 12103        if (_sendStreams.Count > 0)
 2104        {
 12105            _sendStreams.AddLast(message);
 12106            return OutgoingAsyncBase.AsyncStatusQueued;
 2107        }
 2108
 2109        // Prepare the message for sending.
 2110        Debug.Assert(!message.prepared);
 2111
 12112        OutputStream stream = message.stream;
 2113
 12114        message.stream = doCompress(stream, message.compress);
 12115        message.stream.prepareWrite();
 12116        message.prepared = true;
 2117
 12118        TraceUtil.traceSend(stream, _instance, this, _logger, _traceLevels);
 2119
 2120        // Send the message without blocking.
 12121        if (_observer is not null)
 2122        {
 12123            observerStartWrite(message.stream.getBuffer());
 2124        }
 12125        int op = write(message.stream.getBuffer());
 12126        if (op == 0)
 2127        {
 2128            // The message was sent so we're done.
 2129
 12130            if (_observer is not null)
 2131            {
 12132                observerFinishWrite(message.stream.getBuffer());
 2133            }
 2134
 12135            int status = OutgoingAsyncBase.AsyncStatusSent;
 12136            if (message.sent())
 2137            {
 2138                // If there's a sent callback, indicate the caller that it should invoke the sent callback.
 12139                status |= OutgoingAsyncBase.AsyncStatusInvokeSentCallback;
 2140            }
 2141
 12142            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
 12150        _writeStream.swap(message.stream);
 12151        _sendStreams.AddLast(message);
 12152        _threadPool.register(this, op);
 12153        return OutgoingAsyncBase.AsyncStatusQueued;
 2154    }
 2155
 2156    private OutputStream doCompress(OutputStream decompressed, bool compress)
 2157    {
 12158        if (BZip2.isLoaded(_logger) && compress && decompressed.size() >= 100)
 2159        {
 2160            //
 2161            // Do compression.
 2162            //
 12163            Ice.Internal.Buffer cbuf = BZip2.compress(
 12164                decompressed.getBuffer(),
 12165                Protocol.headerSize,
 12166                _compressionLevel);
 12167            if (cbuf is not null)
 2168            {
 12169                var cstream = new OutputStream(new Internal.Buffer(cbuf, true), decompressed.getEncoding());
 2170
 2171                //
 2172                // Set compression status.
 2173                //
 12174                cstream.pos(9);
 12175                cstream.writeByte(2);
 2176
 2177                //
 2178                // Write the size of the compressed stream into the header.
 2179                //
 12180                cstream.pos(10);
 12181                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                //
 12187                decompressed.pos(9);
 12188                decompressed.writeByte(2);
 12189                decompressed.writeInt(cstream.size());
 2190
 12191                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.
 12197        decompressed.pos(9);
 12198        decompressed.writeByte((byte)((BZip2.isLoaded(_logger) && compress) ? 1 : 0));
 2199
 2200        //
 2201        // Not compressed, fill in the message size.
 2202        //
 12203        decompressed.pos(10);
 12204        decompressed.writeInt(decompressed.size());
 2205
 12206        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
 12224        info.stream = new InputStream(_instance, Protocol.currentProtocolEncoding);
 12225        _readStream.swap(info.stream);
 12226        _readStream.resize(Protocol.headerSize);
 12227        _readStream.pos(0);
 12228        _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            //
 12237            info.stream.pos(8);
 12238            byte messageType = info.stream.readByte();
 12239            info.compress = info.stream.readByte();
 12240            if (info.compress == 2)
 2241            {
 12242                if (BZip2.isLoaded(_logger))
 2243                {
 12244                    Ice.Internal.Buffer ubuf = BZip2.decompress(
 12245                        info.stream.getBuffer(),
 12246                        Protocol.headerSize,
 12247                        _messageSizeMax);
 12248                    info.stream = new InputStream(info.stream.instance(), info.stream.getEncoding(), ubuf, true);
 2249                }
 2250                else
 2251                {
 02252                    throw new FeatureNotSupportedException(
 02253                        "Cannot decompress compressed message: BZip2 library is not loaded.");
 2254                }
 2255            }
 12256            info.stream.pos(Protocol.headerSize);
 2257
 2258            switch (messageType)
 2259            {
 2260                case Protocol.closeConnectionMsg:
 2261                {
 12262                    TraceUtil.traceRecv(info.stream, this, _logger, _traceLevels);
 12263                    if (_endpoint.datagram())
 2264                    {
 02265                        if (_warn)
 2266                        {
 02267                            _logger.warning("ignoring close connection message for datagram connection:\n" + _desc);
 2268                        }
 2269                    }
 2270                    else
 2271                    {
 12272                        setState(StateClosingPending, new CloseConnectionException());
 2273
 2274                        //
 2275                        // Notify the transceiver of the graceful connection closure.
 2276                        //
 12277                        int op = _transceiver.closing(false, _exception);
 12278                        if (op != 0)
 2279                        {
 12280                            scheduleCloseTimer();
 12281                            return op;
 2282                        }
 12283                        setState(StateClosed);
 2284                    }
 12285                    break;
 2286                }
 2287
 2288                case Protocol.requestMsg:
 2289                {
 12290                    if (_state >= StateClosing)
 2291                    {
 12292                        TraceUtil.trace(
 12293                            "received request during closing\n(ignored by server, client will retry)",
 12294                            info.stream,
 12295                            this,
 12296                            _logger,
 12297                            _traceLevels);
 2298                    }
 2299                    else
 2300                    {
 12301                        TraceUtil.traceRecv(info.stream, this, _logger, _traceLevels);
 12302                        info.requestId = info.stream.readInt();
 12303                        info.requestCount = 1;
 12304                        info.adapter = _adapter;
 12305                        ++info.upcallCount;
 2306
 12307                        cancelInactivityTimer();
 12308                        ++_dispatchCount;
 2309                    }
 12310                    break;
 2311                }
 2312
 2313                case Protocol.requestBatchMsg:
 2314                {
 12315                    if (_state >= StateClosing)
 2316                    {
 02317                        TraceUtil.trace(
 02318                            "received batch request during closing\n(ignored by server, client will retry)",
 02319                            info.stream,
 02320                            this,
 02321                            _logger,
 02322                            _traceLevels);
 2323                    }
 2324                    else
 2325                    {
 12326                        TraceUtil.traceRecv(info.stream, this, _logger, _traceLevels);
 12327                        int requestCount = info.stream.readInt();
 12328                        if (requestCount < 0)
 2329                        {
 02330                            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;
 12340                        if (requestCount > (info.stream.size() - info.stream.pos()) / minBatchRequestSize)
 2341                        {
 02342                            throw new MarshalException(
 02343                                $"Received batch request with {requestCount} batches, more than the message can contain.
 2344                        }
 2345
 12346                        info.requestCount = requestCount;
 12347                        info.adapter = _adapter;
 12348                        info.upcallCount += info.requestCount;
 2349
 12350                        cancelInactivityTimer();
 12351                        _dispatchCount += info.requestCount;
 2352                    }
 12353                    break;
 2354                }
 2355
 2356                case Protocol.replyMsg:
 2357                {
 12358                    TraceUtil.traceRecv(info.stream, this, _logger, _traceLevels);
 12359                    info.requestId = info.stream.readInt();
 12360                    if (_asyncRequests.TryGetValue(info.requestId, out info.outAsync))
 2361                    {
 12362                        _asyncRequests.Remove(info.requestId);
 2363
 12364                        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                        //
 12371                        OutgoingMessage message = _sendStreams.Count > 0 ? _sendStreams.First.Value : null;
 12372                        if (message is not null && message.outAsync == info.outAsync)
 2373                        {
 12374                            message.receivedReply = true;
 2375                        }
 12376                        else if (info.outAsync.response())
 2377                        {
 12378                            ++info.upcallCount;
 2379                        }
 2380                        else
 2381                        {
 12382                            info.outAsync = null;
 2383                        }
 12384                        if (_closeRequested && _state < StateClosing && _asyncRequests.Count == 0)
 2385                        {
 12386                            doApplicationClose();
 2387                        }
 2388                    }
 12389                    break;
 2390                }
 2391
 2392                case Protocol.validateConnectionMsg:
 2393                {
 12394                    TraceUtil.traceRecv(info.stream, this, _logger, _traceLevels);
 12395                    break;
 2396                }
 2397
 2398                default:
 2399                {
 02400                    TraceUtil.trace(
 02401                        "received unknown message\n(invalid, closing connection)",
 02402                        info.stream,
 02403                        this,
 02404                        _logger,
 02405                        _traceLevels);
 2406
 02407                    throw new ProtocolException($"Received Ice protocol message with unknown type: {messageType}");
 2408                }
 2409            }
 12410        }
 12411        catch (LocalException ex)
 2412        {
 12413            if (_endpoint.datagram())
 2414            {
 02415                if (_warn)
 2416                {
 02417                    _logger.warning("datagram connection exception:\n" + ex.ToString() + "\n" + _desc);
 2418                }
 2419            }
 2420            else
 2421            {
 12422                setState(StateClosed, ex);
 2423            }
 12424        }
 2425
 12426        if (_state == StateHolding)
 2427        {
 2428            // Don't continue reading if the connection is in the holding state.
 02429            return SocketOperation.None;
 2430        }
 12431        else if (_maxDispatches > 0 && _dispatchCount >= _maxDispatches)
 2432        {
 2433            // Don't continue reading if the _maxDispatches limit is reached or exceeded.
 12434            _idleTimeoutTransceiver?.disableIdleCheck();
 12435            return SocketOperation.None;
 2436        }
 2437        else
 2438        {
 2439            // Continue reading.
 12440            return SocketOperation.Read;
 2441        }
 12442    }
 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
 12454        Object dispatcher = adapter?.dispatchPipeline;
 2455
 2456        try
 2457        {
 12458            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.
 12463                var request = new IncomingRequest(requestId, this, adapter, stream);
 2464
 12465                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.
 12470                    _ = dispatchAsync(request);
 2471                }
 2472                else
 2473                {
 2474                    // Received request on a connection without an object adapter.
 12475                    sendResponse(
 12476                        request.current.createOutgoingResponse(new ObjectNotExistException()),
 12477                        isTwoWay: !_endpoint.datagram() && requestId != 0,
 12478                        compress: 0);
 2479                }
 12480                --requestCount;
 2481            }
 2482
 12483            stream.clear();
 12484        }
 02485        catch (LocalException ex) // TODO: catch all exceptions
 2486        {
 2487            // Typically, the IncomingRequest constructor throws an exception, and we can't continue.
 02488            dispatchException(ex, requestCount);
 02489        }
 2490
 2491        async Task dispatchAsync(IncomingRequest request)
 2492        {
 2493            try
 2494            {
 2495                OutgoingResponse response;
 2496
 2497                try
 2498                {
 12499                    response = await dispatcher.dispatchAsync(request).ConfigureAwait(false);
 12500                }
 12501                catch (System.Exception ex)
 2502                {
 12503                    response = request.current.createOutgoingResponse(ex);
 12504                }
 2505
 12506                sendResponse(response, isTwoWay: !_endpoint.datagram() && requestId != 0, compress);
 12507            }
 02508            catch (LocalException ex) // TODO: catch all exceptions to avoid UnobservedTaskException
 2509            {
 02510                dispatchException(ex, requestCount: 1);
 02511            }
 12512        }
 12513    }
 2514
 2515    private void sendResponse(OutgoingResponse response, bool isTwoWay, byte compress)
 2516    {
 12517        bool finished = false;
 2518        try
 2519        {
 12520            lock (_mutex)
 2521            {
 2522                Debug.Assert(_state > StateNotValidated);
 2523
 2524                try
 2525                {
 12526                    if (--_upcallCount == 0)
 2527                    {
 12528                        if (_state == StateFinished)
 2529                        {
 12530                            finished = true;
 12531                            _observer?.detach();
 2532                        }
 12533                        Monitor.PulseAll(_mutex);
 2534                    }
 2535
 12536                    if (_state >= StateClosed)
 2537                    {
 2538                        Debug.Assert(_exception is not null);
 12539                        throw _exception;
 2540                    }
 2541
 12542                    if (isTwoWay)
 2543                    {
 12544                        sendMessage(new OutgoingMessage(response.outputStream, compress > 0));
 2545                    }
 2546
 12547                    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.
 12551                        _threadPool.update(this, SocketOperation.None, SocketOperation.Read);
 12552                        _idleTimeoutTransceiver?.enableIdleCheck();
 2553                    }
 2554
 12555                    --_dispatchCount;
 2556
 12557                    if (_state == StateClosing && _upcallCount == 0)
 2558                    {
 12559                        initiateShutdown();
 2560                    }
 12561                }
 12562                catch (LocalException ex)
 2563                {
 12564                    setState(StateClosed, ex);
 12565                }
 2566            }
 2567        }
 2568        finally
 2569        {
 12570            if (finished && _removeFromFactory is not null)
 2571            {
 12572                _removeFromFactory(this);
 2573            }
 12574        }
 12575    }
 2576
 2577    private void dispatchException(LocalException ex, int requestCount)
 2578    {
 02579        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.
 02583        lock (_mutex)
 2584        {
 02585            setState(StateClosed, ex);
 2586
 02587            if (requestCount > 0)
 2588            {
 2589                Debug.Assert(_upcallCount >= requestCount);
 02590                _upcallCount -= requestCount;
 02591                if (_upcallCount == 0)
 2592                {
 02593                    if (_state == StateFinished)
 2594                    {
 02595                        finished = true;
 02596                        _observer?.detach();
 2597                    }
 02598                    Monitor.PulseAll(_mutex);
 2599                }
 2600            }
 02601        }
 2602
 02603        if (finished && _removeFromFactory is not null)
 2604        {
 02605            _removeFromFactory(this);
 2606        }
 02607    }
 2608
 2609    private void inactivityCheck(System.Threading.Timer inactivityTimer)
 2610    {
 12611        lock (_mutex)
 2612        {
 2613            // If the timers are different, it means this inactivityTimer is no longer current.
 12614            if (inactivityTimer == _inactivityTimer)
 2615            {
 12616                _inactivityTimer = null;
 12617                inactivityTimer.Dispose(); // non-blocking
 2618
 12619                if (_state == StateActive)
 2620                {
 12621                    setState(
 12622                        StateClosing,
 12623                        new ConnectionClosedException(
 12624                            "Connection closed because it remained inactive for longer than the inactivity timeout.",
 12625                            closedByApplication: false));
 2626                }
 2627            }
 2628            // Else this timer was already canceled and disposed. Nothing to do.
 12629        }
 12630    }
 2631
 2632    private void connectTimedOut(System.Threading.Timer connectTimer)
 2633    {
 12634        lock (_mutex)
 2635        {
 12636            if (_state < StateActive)
 2637            {
 12638                setState(StateClosed, new ConnectTimeoutException());
 2639            }
 12640        }
 2641        // else ignore since we're already connected.
 12642        connectTimer.Dispose();
 12643    }
 2644
 2645    private void closeTimedOut(System.Threading.Timer closeTimer)
 2646    {
 12647        lock (_mutex)
 2648        {
 12649            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.
 12653                _exception = new CloseTimeoutException();
 12654                setState(StateClosed);
 2655            }
 12656        }
 2657        // else ignore since we're already closed.
 12658        closeTimer.Dispose();
 12659    }
 2660
 2661    private ConnectionInfo initConnectionInfo()
 2662    {
 2663        // Called with _mutex locked.
 2664
 12665        if (_state > StateNotInitialized && _info is not null) // Update the connection info until it's initialized
 2666        {
 12667            return _info;
 2668        }
 2669
 12670        _info =
 12671            _transceiver.getInfo(incoming: _connector is null, _adapter?.getName() ?? "", _endpoint.connectionId());
 12672        return _info;
 2673    }
 2674
 02675    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    {
 12679        if (_readStreamPos >= 0)
 2680        {
 2681            Debug.Assert(!buf.empty());
 12682            _observer.receivedBytes(buf.b.position() - _readStreamPos);
 2683        }
 12684        _readStreamPos = buf.empty() ? -1 : buf.b.position();
 12685    }
 2686
 2687    private void observerFinishRead(Ice.Internal.Buffer buf)
 2688    {
 12689        if (_readStreamPos == -1)
 2690        {
 02691            return;
 2692        }
 2693        Debug.Assert(buf.b.position() >= _readStreamPos);
 12694        _observer.receivedBytes(buf.b.position() - _readStreamPos);
 12695        _readStreamPos = -1;
 12696    }
 2697
 2698    private void observerStartWrite(Ice.Internal.Buffer buf)
 2699    {
 12700        if (_writeStreamPos >= 0)
 2701        {
 2702            Debug.Assert(!buf.empty());
 12703            _observer.sentBytes(buf.b.position() - _writeStreamPos);
 2704        }
 12705        _writeStreamPos = buf.empty() ? -1 : buf.b.position();
 12706    }
 2707
 2708    private void observerFinishWrite(Ice.Internal.Buffer buf)
 2709    {
 12710        if (_writeStreamPos == -1)
 2711        {
 12712            return;
 2713        }
 12714        if (buf.b.position() > _writeStreamPos)
 2715        {
 12716            _observer.sentBytes(buf.b.position() - _writeStreamPos);
 2717        }
 12718        _writeStreamPos = -1;
 12719    }
 2720
 2721    private int read(Ice.Internal.Buffer buf)
 2722    {
 12723        int start = buf.b.position();
 12724        int op = _transceiver.read(buf, ref _hasMoreData);
 12725        if (_instance.traceLevels().network >= 3 && buf.b.position() != start)
 2726        {
 12727            var s = new StringBuilder("received ");
 12728            if (_endpoint.datagram())
 2729            {
 02730                s.Append(buf.b.limit());
 2731            }
 2732            else
 2733            {
 12734                s.Append(buf.b.position() - start);
 12735                s.Append(" of ");
 12736                s.Append(buf.b.limit() - start);
 2737            }
 12738            s.Append(" bytes via ");
 12739            s.Append(_endpoint.protocol());
 12740            s.Append('\n');
 12741            s.Append(ToString());
 12742            _instance.initializationData().logger.trace(_instance.traceLevels().networkCat, s.ToString());
 2743        }
 12744        return op;
 2745    }
 2746
 2747    private int write(Ice.Internal.Buffer buf)
 2748    {
 12749        int start = buf.b.position();
 12750        int op = _transceiver.write(buf);
 12751        if (_instance.traceLevels().network >= 3 && buf.b.position() != start)
 2752        {
 12753            var s = new StringBuilder("sent ");
 12754            s.Append(buf.b.position() - start);
 12755            if (!_endpoint.datagram())
 2756            {
 12757                s.Append(" of ");
 12758                s.Append(buf.b.limit() - start);
 2759            }
 12760            s.Append(" bytes via ");
 12761            s.Append(_endpoint.protocol());
 12762            s.Append('\n');
 12763            s.Append(ToString());
 12764            _instance.initializationData().logger.trace(_instance.traceLevels().networkCat, s.ToString());
 2765        }
 12766        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
 12775        _inactivityTimer = new System.Threading.Timer(
 12776            inactivityTimer => inactivityCheck((System.Threading.Timer)inactivityTimer));
 12777        _inactivityTimer.Change(_inactivityTimeout, Timeout.InfiniteTimeSpan);
 12778    }
 2779
 2780    private void cancelInactivityTimer()
 2781    {
 2782        // Called with the ConnectionI mutex locked.
 12783        if (_inactivityTimer is not null)
 2784        {
 12785            _inactivityTimer.Change(Timeout.InfiniteTimeSpan, Timeout.InfiniteTimeSpan);
 12786            _inactivityTimer.Dispose();
 12787            _inactivityTimer = null;
 2788        }
 12789    }
 2790
 2791    private void scheduleCloseTimer()
 2792    {
 12793        if (_closeTimeout > TimeSpan.Zero)
 2794        {
 2795#pragma warning disable CA2000 // closeTimer is disposed by closeTimedOut.
 12796            var closeTimer = new System.Threading.Timer(
 12797                timerObj => closeTimedOut((System.Threading.Timer)timerObj));
 2798            // schedule timer to run once; closeTimedOut disposes the timer too.
 12799            closeTimer.Change(_closeTimeout, Timeout.InfiniteTimeSpan);
 2800#pragma warning restore CA2000
 2801        }
 12802    }
 2803
 2804    private void doApplicationClose()
 2805    {
 2806        // Called with the ConnectionI mutex locked.
 2807        Debug.Assert(_state < StateClosing);
 12808        setState(
 12809            StateClosing,
 12810            new ConnectionClosedException(
 12811                "The connection was closed gracefully by the application.",
 12812                closedByApplication: true));
 12813    }
 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
 12870    private static readonly ConnectionState[] connectionStateMap = [
 12871        ConnectionState.ConnectionStateValidating,   // StateNotInitialized
 12872        ConnectionState.ConnectionStateValidating,   // StateNotValidated
 12873        ConnectionState.ConnectionStateActive,       // StateActive
 12874        ConnectionState.ConnectionStateHolding,      // StateHolding
 12875        ConnectionState.ConnectionStateClosing,      // StateClosing
 12876        ConnectionState.ConnectionStateClosing,      // StateClosingPending
 12877        ConnectionState.ConnectionStateClosed,       // StateClosed
 12878        ConnectionState.ConnectionStateClosed,       // StateFinished
 12879    ];
 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
 12915    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
 12922    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.
 12968    private readonly TaskCompletionSource _closed = new(TaskCreationOptions.RunContinuationsAsynchronously);
 12969    private readonly object _mutex = new();
 2970}

Methods/Properties

start(Ice.ConnectionI.StartCallback)
startAndWait()
activate()
hold()
destroy(int)
abort()
closeAsync()
isActiveOrHolding()
throwException()
waitUntilHolding()
waitUntilFinished()
updateObserver()
sendAsyncRequest(Ice.Internal.OutgoingAsyncBase, bool, bool, int)
getBatchRequestQueue()
flushBatchRequests(Ice.CompressBatch)
flushBatchRequestsAsync(Ice.CompressBatch, System.IProgress<bool>, System.Threading.CancellationToken)
disableInactivityCheck()
setCloseCallback(Ice.CloseCallback)
asyncRequestCanceled(Ice.Internal.OutgoingAsyncBase, Ice.LocalException)
endpoint()
connector()
setAdapter(Ice.ObjectAdapter)
getAdapter()
getEndpoint()
createProxy(Ice.Identity)
setAdapterFromAdapter(Ice.ObjectAdapter)
startAsync(int, Ice.Internal.AsyncCallback)
doIO()
finishAsync(int)
message(Ice.Internal.ThreadPoolCurrent)
upcall(Ice.ConnectionI.StartCallback, System.Collections.Generic.Queue<Ice.ConnectionI.OutgoingMessage>, Ice.ConnectionI.MessageInfo)
finished(Ice.Internal.ThreadPoolCurrent)
finish()
ToString()
type()
getInfo()
setBufferSize(int, int)
exception(Ice.LocalException)
getThreadPool()
.ctor(Ice.Internal.Instance, Ice.Internal.Transceiver, Ice.Internal.Connector, Ice.Internal.EndpointI, Ice.ObjectAdapter, System.Action<Ice.ConnectionI>, Ice.ConnectionOptions)
idleCheck(System.TimeSpan)
sendHeartbeat()
isHeartbeat()
toConnectionState(int)
setState(int, Ice.LocalException)
setState(int)
initiateShutdown()
initialize(int)
validate(int)
sendNextMessage(out System.Collections.Generic.Queue<Ice.ConnectionI.OutgoingMessage>)
sendMessage(Ice.ConnectionI.OutgoingMessage)
doCompress(Ice.OutputStream, bool)
parseMessage(ref Ice.ConnectionI.MessageInfo)
dispatchAll(Ice.InputStream, int, int, byte, Ice.ObjectAdapter)
dispatchAsync()
sendResponse(Ice.OutgoingResponse, bool, byte)
dispatchException(Ice.LocalException, int)
inactivityCheck(System.Threading.Timer)
connectTimedOut(System.Threading.Timer)
closeTimedOut(System.Threading.Timer)
initConnectionInfo()
warning(string, System.Exception)
observerStartRead(Ice.Internal.Buffer)
observerFinishRead(Ice.Internal.Buffer)
observerStartWrite(Ice.Internal.Buffer)
observerFinishWrite(Ice.Internal.Buffer)
read(Ice.Internal.Buffer)
write(Ice.Internal.Buffer)
scheduleInactivityTimer()
cancelInactivityTimer()
scheduleCloseTimer()
doApplicationClose()
.ctor(Ice.OutputStream, bool)
.ctor(Ice.Internal.OutgoingAsyncBase, Ice.OutputStream, bool, int)
canceled()
sent()
completed(Ice.LocalException)
.cctor()