< Summary

Information
Class: Ice.Internal.BatchRequestQueue
Assembly: Ice
File(s): /_/csharp/src/Ice/Internal/BatchRequestQueue.cs
Tag: 125_37167941578
Line coverage
88%
Covered lines: 83
Uncovered lines: 11
Coverable lines: 94
Total lines: 231
Line coverage: 88.2%
Branch coverage
83%
Covered branches: 25
Total branches: 30
Branch coverage: 83.3%
Method coverage
85%
Covered methods: 6
Fully covered methods: 3
Total methods: 7
Method coverage: 85.7%
Full method coverage: 42.8%

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
.ctor(...)100%44100%
prepareBatchRequest(...)50%4477.78%
finishBatchRequest(...)100%88100%
abortBatchRequest(...)0%620%
swap(...)90%101096%
destroy(...)100%11100%
enqueueBatchRequest(...)100%22100%

File(s)

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

#LineLine coverage
 1// Copyright (c) ZeroC, Inc.
 2
 3using System.Diagnostics;
 4
 5namespace Ice.Internal;
 6
 7internal sealed class BatchRequestI : BatchRequest
 8{
 9    internal BatchRequestI(BatchRequestQueue queue) => _queue = queue;
 10
 11    internal void reset(ObjectPrx proxy, string operation, int size)
 12    {
 13        _proxy = proxy;
 14        _operation = operation;
 15        _size = size;
 16    }
 17
 18    public void enqueue() => _queue.enqueueBatchRequest(_proxy);
 19
 20    public ObjectPrx getProxy() => _proxy;
 21
 22    public string getOperation() => _operation;
 23
 24    public int getSize() => _size;
 25
 26    private readonly BatchRequestQueue _queue;
 27    private ObjectPrx _proxy;
 28    private string _operation;
 29    private int _size;
 30}
 31
 32internal sealed class BatchRequestQueue
 33{
 134    internal BatchRequestQueue(Instance instance, bool datagram)
 35    {
 136        InitializationData initData = instance.initializationData();
 137        _interceptor = initData.batchRequestInterceptor;
 138        _batchStreamInUse = false;
 139        _batchRequestNum = 0;
 140        _batchStream =
 141            new OutputStream(Protocol.currentProtocolEncoding, instance.defaultsAndOverrides().defaultFormat);
 142        _batchStream.writeBlob(Protocol.requestBatchHdr);
 143        _batchMarker = _batchStream.size();
 144        _request = new BatchRequestI(this);
 45
 146        _maxSize = instance.batchAutoFlushSize();
 147        if (_maxSize > 0 && datagram)
 48        {
 149            int udpSndSize = initData.properties.getPropertyAsIntWithDefault(
 150                "Ice.UDP.SndSize",
 151                65535 - _udpOverhead);
 152            if (udpSndSize < _maxSize)
 53            {
 154                _maxSize = udpSndSize;
 55            }
 56        }
 157    }
 58
 59    internal void prepareBatchRequest(OutputStream os)
 60    {
 161        lock (_mutex)
 62        {
 163            if (_exception != null)
 64            {
 065                throw _exception;
 66            }
 67
 68            // This is similar to a mutex lock in that the stream is only "locked" while marshaling.
 169            while (_batchStreamInUse)
 70            {
 071                Monitor.Wait(_mutex);
 72            }
 173            _batchStreamInUse = true;
 174            _batchStream.swap(os);
 175        }
 176    }
 77
 78    internal void finishBatchRequest(OutputStream os, ObjectPrx proxy, string operation)
 79    {
 180        lock (_mutex)
 81        {
 82            // Bring the request stream back and become the owner so this thread's own auto-flush below can
 83            // re-enter swap(). swap() lets only _batchStreamOwner through while the stream is in use, and only
 84            // after passing that gate does it read the queue's bookkeeping; no other thread can touch the
 85            // queue, so the mutations below need no further synchronization.
 86            Debug.Assert(_batchStreamInUse);
 187            _batchStream.swap(os);
 188            _batchStreamOwner = Thread.CurrentThread;
 189        }
 90
 91        try
 92        {
 193            if (_maxSize > 0 && _batchStream.size() >= _maxSize)
 94            {
 195                _ = proxy.ice_flushBatchRequestsAsync(); // Auto flush
 96            }
 97
 98            Debug.Assert(_batchMarker < _batchStream.size());
 199            if (_interceptor != null)
 100            {
 1101                _request.reset(proxy, operation, _batchStream.size() - _batchMarker);
 1102                _interceptor(_request, _batchRequestNum, _batchMarker);
 103            }
 104            else
 105            {
 1106                bool? compress = ((ObjectPrxHelperBase)proxy).iceReference().getCompressOverride();
 1107                if (compress is not null)
 108                {
 1109                    _batchCompress |= compress.Value;
 110                }
 1111                _batchMarker = _batchStream.size();
 1112                ++_batchRequestNum;
 113            }
 1114        }
 115        finally
 116        {
 1117            lock (_mutex)
 118            {
 1119                _batchStream.resize(_batchMarker);
 1120                _batchStreamInUse = false;
 1121                _batchStreamOwner = null;
 1122                Monitor.PulseAll(_mutex);
 1123            }
 1124        }
 1125    }
 126
 127    internal void abortBatchRequest(OutputStream os)
 128    {
 0129        lock (_mutex)
 130        {
 0131            if (_batchStreamInUse)
 132            {
 0133                _batchStream.swap(os);
 0134                _batchStream.resize(_batchMarker);
 0135                _batchStreamInUse = false;
 0136                Monitor.PulseAll(_mutex);
 137            }
 0138        }
 0139    }
 140
 141    internal int swap(OutputStream os, out bool compress)
 142    {
 1143        lock (_mutex)
 144        {
 145            // Only the thread that currently owns the in-use stream may bypass the wait, to run its own
 146            // auto-flush; every other caller waits for the in-progress batch request to finish.
 1147            if (_batchStreamOwner != Thread.CurrentThread)
 148            {
 1149                while (_batchStreamInUse)
 150                {
 0151                    Monitor.Wait(_mutex);
 152                }
 153            }
 154
 155            // Read the bookkeeping only after passing the gate above: finishBatchRequest mutates these scalars
 156            // without the lock while it owns the in-use stream, so reading _batchRequestNum before the wait
 157            // would race those writes.
 1158            if (_batchRequestNum == 0)
 159            {
 1160                compress = false;
 1161                return 0;
 162            }
 163
 1164            byte[] lastRequest = null;
 1165            if (_batchMarker < _batchStream.size())
 166            {
 1167                lastRequest = new byte[_batchStream.size() - _batchMarker];
 1168                Buffer buffer = _batchStream.getBuffer();
 1169                buffer.b.position(_batchMarker);
 1170                buffer.b.get(lastRequest);
 1171                _batchStream.resize(_batchMarker);
 172            }
 173
 1174            int requestNum = _batchRequestNum;
 1175            compress = _batchCompress;
 1176            _batchStream.swap(os);
 177
 178            //
 179            // Reset the batch.
 180            //
 1181            _batchRequestNum = 0;
 1182            _batchCompress = false;
 1183            _batchStream.writeBlob(Protocol.requestBatchHdr);
 1184            _batchMarker = _batchStream.size();
 1185            if (lastRequest != null)
 186            {
 1187                _batchStream.writeBlob(lastRequest);
 188            }
 1189            return requestNum;
 190        }
 1191    }
 192
 193    internal void destroy(LocalException ex)
 194    {
 1195        lock (_mutex)
 196        {
 1197            _exception = ex;
 1198        }
 1199    }
 200
 201    internal void enqueueBatchRequest(ObjectPrx proxy)
 202    {
 203        Debug.Assert(_batchMarker < _batchStream.size());
 1204        bool? compress = ((ObjectPrxHelperBase)proxy).iceReference().getCompressOverride();
 1205        if (compress is not null)
 206        {
 1207            _batchCompress |= compress.Value;
 208        }
 1209        _batchMarker = _batchStream.size();
 1210        ++_batchRequestNum;
 1211    }
 212
 1213    private readonly object _mutex = new();
 214
 215    private readonly System.Action<BatchRequest, int, int> _interceptor;
 216    private readonly OutputStream _batchStream;
 217    private bool _batchStreamInUse;
 218
 219    // While finishBatchRequest holds the in-use stream, this is its thread: the only thread allowed to
 220    // re-enter swap() (for its own auto-flush). A null owner means no thread may flush the in-use stream, so
 221    // other threads wait for _batchStreamInUse to clear.
 222    private Thread _batchStreamOwner;
 223
 224    private int _batchRequestNum;
 225    private int _batchMarker;
 226    private bool _batchCompress;
 227    private readonly BatchRequestI _request;
 228    private LocalException _exception;
 229    private readonly int _maxSize;
 230    private const int _udpOverhead = 20 + 8;
 231}