Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
91 changes: 46 additions & 45 deletions Source/FileWatcherEx/Helpers/EventProcessor.cs
Original file line number Diff line number Diff line change
Expand Up @@ -12,9 +12,9 @@ internal class EventProcessor
private const int EventDelay = 50;

/// <summary>
/// Warn after certain time span of event spam (in ticks)
/// Warn after certain time span of event spam
/// </summary>
private const int EventSpamWarningThreshold = 60 * 1000 * 10000;
private readonly TimeSpan _eventSpamWarningThreshold = TimeSpan.FromMinutes(1);

private readonly object _lock = new();
private Task? _delayTask = null;
Expand Down Expand Up @@ -42,20 +42,7 @@ public void ProcessEvent(FileChangedEvent fileEvent)
lock (_lock)
{
var now = DateTime.Now.Ticks;

// Check for spam
if (_events.Count == 0)
{
_spamWarningLogged = false;
_spamCheckStartTime = now;
}
else if (!_spamWarningLogged && _spamCheckStartTime + EventSpamWarningThreshold < now)
{
_spamWarningLogged = true;
_logger(string.Format(
"Warning: Watcher is busy catching up with {0} file changes in 60 seconds. Latest path is '{1}'",
_events.Count, fileEvent.FullPath));
}
WarnForSpam(fileEvent, now);

// Add into our queue
_events.Add(fileEvent);
Expand All @@ -64,40 +51,54 @@ public void ProcessEvent(FileChangedEvent fileEvent)
// Process queue after delay
if (_delayTask == null)
{
// Create function to buffer events
void Func(Task value)
// Start function after delay
_delayStarted = _lastEventTime;
_delayTask = Task.Delay(EventDelay).ContinueWith(HandleEventsFunc);
}
}
}

private void HandleEventsFunc(Task _)
{
lock (_lock)
{
// Check if another event has been received in the meantime
if (_delayStarted == _lastEventTime)
{
// Normalize and handle
var normalized = new EventNormalizer().Normalize(_events.ToArray());
foreach (var ev in normalized)
{
lock (_lock)
{
// Check if another event has been received in the meantime
if (_delayStarted == _lastEventTime)
{
// Normalize and handle
var normalized = new EventNormalizer().Normalize(_events.ToArray());
foreach (var e in normalized)
{
_handleEvent(e);
}

// Reset
_events.Clear();
_delayTask = null;
}

// Otherwise we have received a new event while this task was
// delayed and we reschedule it.
else
{
_delayStarted = _lastEventTime;
_delayTask = Task.Delay(EventDelay).ContinueWith(Func);
}
}
_handleEvent(ev);
}

// Start function after delay
// Reset
_events.Clear();
_delayTask = null;
}

// Otherwise we have received a new event while this task was
// delayed and we reschedule it.
else
{
_delayStarted = _lastEventTime;
_delayTask = Task.Delay(EventDelay).ContinueWith(Func);
_delayTask = Task.Delay(EventDelay).ContinueWith(HandleEventsFunc);
}
}
}

private void WarnForSpam(FileChangedEvent fileEvent, long now)
{
if (_events.Count == 0)
{
_spamWarningLogged = false;
_spamCheckStartTime = now;
}
else if (! _spamWarningLogged && _spamCheckStartTime + _eventSpamWarningThreshold.Ticks < now)
{
_spamWarningLogged = true;
_logger($"Warning: Watcher is busy catching up with {_events.Count} file changes " +
$"in {_eventSpamWarningThreshold.TotalSeconds} seconds. Latest path is '{fileEvent.FullPath}'");
}
}
}