diff --git a/Duplicati/Library/Utility/ThrottledStream.cs b/Duplicati/Library/Utility/ThrottledStream.cs index c0a7afe14..37b8f940f 100644 --- a/Duplicati/Library/Utility/ThrottledStream.cs +++ b/Duplicati/Library/Utility/ThrottledStream.cs @@ -1,263 +1,261 @@ -#region Disclaimer / License -// Copyright (C) 2015, The Duplicati Team -// http://www.duplicati.com, info@duplicati.com -// -// This library is free software; you can redistribute it and/or -// modify it under the terms of the GNU Lesser General Public -// License as published by the Free Software Foundation; either -// version 2.1 of the License, or (at your option) any later version. -// -// This library is distributed in the hope that it will be useful, -// but WITHOUT ANY WARRANTY; without even the implied warranty of -// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU -// Lesser General Public License for more details. -// -// You should have received a copy of the GNU Lesser General Public -// License along with this library; if not, write to the Free Software -// Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA -// -#endregion -using System; -using System.Collections.Generic; -using System.Text; -using System.IO; - -//TODO: Use the IPGlobalProperties to dynamically throttle data -//http://msdn.microsoft.com/en-us/library/system.net.networkinformation.ipglobalproperties.aspx - -namespace Duplicati.Library.Utility -{ - /// - /// This class throttles the rate data can be read or written to the underlying stream. - /// This creates a bandwith throttle option for any stream, including a network stream. - /// - public class ThrottledStream : OverrideableStream - { - /// - /// The delegate type for the callback - /// - public delegate void ThrottledStreamCallback(ThrottledStream sender); - - /// - /// An event that is raised while the stream is active - /// - public event ThrottledStreamCallback Callback; - - /// - /// The max number of bytes pr. second to write - /// - private long m_writespeed; - /// - /// The max number of bytes pr. second to read - /// - private long m_readspeed; - - /// - /// This is a list of the most recent reads. The key is the tick at the time, and the value is the number of bytes. - /// - List> m_dataread; - /// - /// This is a list of the most recent writes. The key is the tick at the time, and the value is the number of bytes. - /// - List> m_datawritten; - /// - /// This is the sum of all bytes in the m_dataread table, summed for fast access. - /// - private long m_bytesread; - /// - /// This is the sum of all bytes in the m_datawritten table, summed for fast access. - /// - private long m_byteswritten; - - /// - /// The number of bytes transfered without raising an event - /// - private long m_progresscounter = 0; - - /// - /// The number of reads or writes to keep track of - /// - private const long STATISTICS_SIZE = 500; - /// - /// The number of ticks to have passed before the throttle begins - /// - private const long MIN_DURATION = TimeSpan.TicksPerSecond / 4; - /// - /// The number of sub chunks to perform when throttling - /// - private const int DELAY_CHUNK_SIZE = 1024; - - /// - /// The number of bytes to process without raising an event - /// - private const int REPORT_DISTANCE_SIZE = 1024 * 50; - - /// - /// Creates a throttle around a stream. - /// - /// The stream to throttle - /// The maximum number of bytes pr. second to read. Specify a number less than 1 to allow unlimited speed. - /// The maximum number of bytes pr. second to write. Specify a number less than 1 to allow unlimited speed. - public ThrottledStream(Stream basestream, long readspeed, long writespeed) - : base(basestream) - { - m_readspeed = readspeed; - m_writespeed = writespeed; - - if (m_basestream.CanRead && m_readspeed > 0) - m_dataread = new List>(); - if (m_basestream.CanWrite && m_writespeed > 0) - m_datawritten = new List>(); - } - - public override int Read(byte[] buffer, int offset, int count) - { - int bytesRead = DelayIfRequired(true, buffer, ref offset, ref count); - if (count != 0) - bytesRead += m_basestream.Read(buffer, offset, count); - - m_progresscounter += bytesRead; - - if (m_progresscounter > REPORT_DISTANCE_SIZE) - { - m_progresscounter %= REPORT_DISTANCE_SIZE; - if (Callback != null) - Callback(this); - } - - return bytesRead; - } - - public override void Write(byte[] buffer, int offset, int count) - { - m_progresscounter += count; - - DelayIfRequired(false, buffer, ref offset, ref count); - if (count > 0) - m_basestream.Write(buffer, offset, count); - - if (m_progresscounter > REPORT_DISTANCE_SIZE) - { - m_progresscounter %= REPORT_DISTANCE_SIZE; - if (Callback != null) - Callback(this); - } - } - - /// - /// Gets or sets the current read speed in bytes pr. second. - /// Set to zero or less to disable throttling. - /// - public long ReadSpeed - { - get { return m_readspeed; } - set { m_readspeed = value; } - } - - /// - /// Gets or sets the current write speed in bytes pr. second. - /// Set to zero or less to disable throttling - /// - public long WriteSpeed - { - get { return m_writespeed; } - set { m_writespeed = value; } - } - - /// - /// Calculates the speed, and inserts appropriate delays - /// - /// True if the operation is read, false otherwise - /// The data buffer - /// The offset into the buffer - /// The number of bytes to read or write - /// The number of bytes processed while delaying - private int DelayIfRequired(bool read, byte[] buffer, ref int offset, ref int count) - { - if (count <= 0) - return 0; - - List> table = read ? m_dataread : m_datawritten; - int bytesprocessed = 0; - - if (table != null) - { - long maxspeed = read ? m_readspeed : m_writespeed; - Stream stream = m_basestream; - long bytecount = read ? m_bytesread : m_byteswritten; - - //Add this access - table.Add(new KeyValuePair(DateTime.Now.Ticks, count)); - bytecount += count; - - //Prevent too large tables - while (table.Count > STATISTICS_SIZE) - { - bytecount -= table[0].Value; - table.RemoveAt(0); - } - - if (table.Count != 0 && bytecount != 0) - { - TimeSpan duration = new TimeSpan(table[table.Count - 1].Key - table[0].Key); - - //Bail if we are too slow - if (duration.Ticks < MIN_DURATION || bytecount <= 0) - return 0; - - //TODO: The resolution is too low in "TotalSeconds", so the speed is a little higher - double speed = bytecount / duration.TotalSeconds; - if (speed > maxspeed) - { - //We are too fast, delay the access. Calculating how much wait we need. - double secondsNeeded = (bytecount / (double)maxspeed) - duration.TotalSeconds; - long delayTicks = (long)(secondsNeeded * TimeSpan.TicksPerSecond); - - //Calculate the time we should finish, to obey the limit - long targetTime = DateTime.Now.Ticks + delayTicks; - - if (delayTicks > 0) - { - while (count > DELAY_CHUNK_SIZE && delayTicks > 0) - { - int bytes = read ? stream.Read(buffer, offset, DELAY_CHUNK_SIZE) : DELAY_CHUNK_SIZE; - if (!read) - stream.Write(buffer, offset, DELAY_CHUNK_SIZE); - - delayTicks = targetTime - DateTime.Now.Ticks; - long ticksToDelay = (delayTicks / count) * bytes; - - if (ticksToDelay > 0) - System.Threading.Thread.Sleep(new TimeSpan(ticksToDelay)); - - //Reset to include the waited time - delayTicks = targetTime - DateTime.Now.Ticks; - - offset += bytes; - count -= bytes; - bytesprocessed += bytes; - - if (bytes == 0) - break; - } - - if (delayTicks > 0) - System.Threading.Thread.Sleep(new TimeSpan(delayTicks)); - - //Add a marker, indicating that we already slowed down - table.Add(new KeyValuePair(DateTime.Now.Ticks, 0)); - } - - } - } - - if (read) - m_bytesread = bytecount; - else - m_byteswritten = bytecount; - } - - return bytesprocessed; - } - } -} +#region Disclaimer / License +// Copyright (C) 2015, The Duplicati Team +// http://www.duplicati.com, info@duplicati.com +// +// This library is free software; you can redistribute it and/or +// modify it under the terms of the GNU Lesser General Public +// License as published by the Free Software Foundation; either +// version 2.1 of the License, or (at your option) any later version. +// +// This library is distributed in the hope that it will be useful, +// but WITHOUT ANY WARRANTY; without even the implied warranty of +// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU +// Lesser General Public License for more details. +// +// You should have received a copy of the GNU Lesser General Public +// License along with this library; if not, write to the Free Software +// Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +// +#endregion +using System; +using System.Collections.Generic; +using System.Text; +using System.IO; + +//TODO: Use the IPGlobalProperties to dynamically throttle data +//http://msdn.microsoft.com/en-us/library/system.net.networkinformation.ipglobalproperties.aspx + +namespace Duplicati.Library.Utility +{ + /// + /// This class throttles the rate data can be read or written to the underlying stream. + /// This creates a bandwith throttle option for any stream, including a network stream. + /// + public class ThrottledStream : OverrideableStream + { + /// + /// The delegate type for the callback + /// + public delegate void ThrottledStreamCallback(ThrottledStream sender); + + /// + /// An event that is raised while the stream is active + /// + public event ThrottledStreamCallback Callback; + + /// + /// The max number of bytes pr. second to write + /// + private long m_writespeed; + /// + /// The max number of bytes pr. second to read + /// + private long m_readspeed; + + /// + /// The time the last read was sampled + /// + private DateTime m_last_read_sample; + /// + /// The bytes-read counter for this period + /// + private long m_current_read_counter; + /// + /// The current measured read speed + /// + private double m_current_read_speed; + + /// + /// The time the last read was sampled + /// + private DateTime m_last_write_sample; + /// + /// The bytes-written counter for this period + /// + private long m_current_write_counter; + /// + /// The current measured read speed + /// + private double m_current_write_speed; + + /// + /// The number of bytes transfered without raising an event + /// + private long m_progresscounter = 0; + + /// + /// The number of ticks to have passed before a sample is taken + /// + private const long SAMPLE_PERIOD = TimeSpan.TicksPerSecond / 4; + + /// + /// The number of bytes to process without raising an event + /// + private const int REPORT_DISTANCE_SIZE = 1024 * 50; + + /// + /// Creates a throttle around a stream. + /// + /// The stream to throttle + /// The maximum number of bytes pr. second to read. Specify a number less than 1 to allow unlimited speed. + /// The maximum number of bytes pr. second to write. Specify a number less than 1 to allow unlimited speed. + public ThrottledStream(Stream basestream, long readspeed, long writespeed) + : base(basestream) + { + m_readspeed = readspeed; + m_writespeed = writespeed; + + m_last_read_sample = m_last_write_sample = new DateTime(0); + } + + /// + /// Read the specified buffer, offset and count. + /// + /// The buffer to read from. + /// The offset into the buffer. + /// The number of bytes to read. + public override int Read(byte[] buffer, int offset, int count) + { + var remaining = count; + + while (remaining > 0) + { + // To avoid excessive waiting, the delay will wait at most 2 seconds, + // so we split the blocks to limit the number of seconds we can wait + var chunksize = (int)Math.Min(remaining, m_readspeed <= 0 ? remaining : m_readspeed * 2); + DelayIfRequired(ref m_readspeed, chunksize, ref m_last_read_sample, ref m_current_read_counter, ref m_current_read_speed); + + var actual = m_basestream.Read(buffer, offset, chunksize); + + if (actual <= 0) + break; + + m_progresscounter += actual; + m_current_read_counter += actual; + + if (m_progresscounter > REPORT_DISTANCE_SIZE) + { + m_progresscounter %= REPORT_DISTANCE_SIZE; + if (Callback != null) + Callback(this); + } + + remaining -= actual; + } + + return count - remaining; + } + + /// + /// Write the specified buffer, offset and count. + /// + /// The buffer to write to. + /// The offset into the buffer. + /// The number of bytes to write. + public override void Write(byte[] buffer, int offset, int count) + { + while (count > 0) + { + // To avoid excessive waiting, the delay will wait at most 2 seconds, + // so we split the blocks to limit the number of seconds we can wait + var chunksize = (int)Math.Min(count, m_writespeed <= 0 ? count : m_writespeed * 2); + DelayIfRequired(ref m_writespeed, chunksize, ref m_last_write_sample, ref m_current_write_counter, ref m_current_write_speed); + m_basestream.Write(buffer, offset, chunksize); + + m_progresscounter += chunksize; + m_current_write_counter += chunksize; + + if (m_progresscounter > REPORT_DISTANCE_SIZE) + { + m_progresscounter %= REPORT_DISTANCE_SIZE; + if (Callback != null) + Callback(this); + } + + count -= chunksize; + } + } + + /// + /// Gets or sets the current read speed in bytes pr. second. + /// Set to zero or less to disable throttling. + /// + public long ReadSpeed + { + get { return m_readspeed; } + set { m_readspeed = value; } + } + + /// + /// Gets or sets the current write speed in bytes pr. second. + /// Set to zero or less to disable throttling + /// + public long WriteSpeed + { + get { return m_writespeed; } + set { m_writespeed = value; } + } + + /// + /// Gets the actual measured read speed. + /// + public double MeasuredReadSpeed { get { return m_current_read_speed; } } + + /// + /// Gets the actual measured write speed. + /// + public double MeasuredWriteSpeed { get { return m_current_write_speed; } } + + /// + /// Calculates the speed, and inserts appropriate delays + /// + /// True if the operation is read, false otherwise + /// The data buffer + /// The offset into the buffer + /// The number of bytes to read or write + /// The number of bytes processed while delaying + private void DelayIfRequired(ref long limit, int count, ref DateTime last_sample, ref long last_count, ref double current_speed) + { + var now = DateTime.Now; + + if (count <= 0 || limit <= 0) + return; + + // If we are just starting, set the timer and counter + if (last_sample.Ticks == 0) + { + last_count = 0; + last_sample = now; + current_speed = limit; + return; + } + + // Compute the current duration + var duration = now - last_sample; + + // Update speed in intervals + if (duration.Ticks > SAMPLE_PERIOD || last_count > limit) + { + // After a sample period, measure how far ahead we are + var target_delay = TimeSpan.FromSeconds(last_count / (double)limit) - duration; + + // If we are actually ahead, delay for a little while + if (target_delay.Ticks > 1000) + { + // With large changes, we avoid sleeping for several minutes + // This makes the throttling more resposive when increasing the + // throughput, even with large changes + var ms = (int)Math.Min(target_delay.TotalMilliseconds, 2 * 1000); + System.Threading.Thread.Sleep(ms); + + // When we compute how fast this sample was, we include the delay + now = DateTime.Now; + } + + current_speed = last_count / (now - last_sample).TotalSeconds; + last_sample = now; + last_count = 0; + } + } + } +}