Fixes formatting (#200)
This commit is contained in:
parent
1f651992d7
commit
86b7b3aed1
209 changed files with 51326 additions and 50397 deletions
|
|
@ -1,80 +1,82 @@
|
|||
/***************************************************************************
|
||||
* BinaryMemoryWriter.cs
|
||||
* -------------------
|
||||
* begin : May 1, 2002
|
||||
* copyright : (C) The RunUO Software Team
|
||||
* email : info@runuo.com
|
||||
*
|
||||
* $Id$
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
/***************************************************************************
|
||||
*
|
||||
* This program is free software; you can redistribute it and/or modify
|
||||
* it under the terms of the GNU General Public License as published by
|
||||
* the Free Software Foundation; either version 2 of the License, or
|
||||
* (at your option) any later version.
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
using System.IO;
|
||||
|
||||
namespace Server
|
||||
{
|
||||
public sealed class BinaryMemoryWriter : BinaryFileWriter
|
||||
{
|
||||
private static byte[] indexBuffer;
|
||||
private readonly MemoryStream stream;
|
||||
|
||||
public BinaryMemoryWriter()
|
||||
: base(new MemoryStream(512), true) =>
|
||||
stream = UnderlyingStream as MemoryStream;
|
||||
|
||||
protected override int BufferSize => 512;
|
||||
|
||||
public int CommitTo(SequentialFileWriterStream dataFile, SequentialFileWriterStream indexFile, int typeCode, uint serial)
|
||||
{
|
||||
Flush();
|
||||
|
||||
var buffer = stream.GetBuffer();
|
||||
var length = (int)stream.Length;
|
||||
|
||||
var position = dataFile.Position;
|
||||
|
||||
dataFile.Write(buffer, 0, length);
|
||||
|
||||
indexBuffer ??= new byte[20];
|
||||
|
||||
indexBuffer[0] = (byte)typeCode;
|
||||
indexBuffer[1] = (byte)(typeCode >> 8);
|
||||
indexBuffer[2] = (byte)(typeCode >> 16);
|
||||
indexBuffer[3] = (byte)(typeCode >> 24);
|
||||
|
||||
indexBuffer[4] = (byte)serial;
|
||||
indexBuffer[5] = (byte)(serial >> 8);
|
||||
indexBuffer[6] = (byte)(serial >> 16);
|
||||
indexBuffer[7] = (byte)(serial >> 24);
|
||||
|
||||
indexBuffer[8] = (byte)position;
|
||||
indexBuffer[9] = (byte)(position >> 8);
|
||||
indexBuffer[10] = (byte)(position >> 16);
|
||||
indexBuffer[11] = (byte)(position >> 24);
|
||||
indexBuffer[12] = (byte)(position >> 32);
|
||||
indexBuffer[13] = (byte)(position >> 40);
|
||||
indexBuffer[14] = (byte)(position >> 48);
|
||||
indexBuffer[15] = (byte)(position >> 56);
|
||||
|
||||
indexBuffer[16] = (byte)length;
|
||||
indexBuffer[17] = (byte)(length >> 8);
|
||||
indexBuffer[18] = (byte)(length >> 16);
|
||||
indexBuffer[19] = (byte)(length >> 24);
|
||||
|
||||
indexFile.Write(indexBuffer, 0, indexBuffer.Length);
|
||||
|
||||
stream.SetLength(0);
|
||||
|
||||
return length;
|
||||
}
|
||||
}
|
||||
}
|
||||
/***************************************************************************
|
||||
* BinaryMemoryWriter.cs
|
||||
* -------------------
|
||||
* begin : May 1, 2002
|
||||
* copyright : (C) The RunUO Software Team
|
||||
* email : info@runuo.com
|
||||
*
|
||||
* $Id$
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
/***************************************************************************
|
||||
*
|
||||
* This program is free software; you can redistribute it and/or modify
|
||||
* it under the terms of the GNU General Public License as published by
|
||||
* the Free Software Foundation; either version 2 of the License, or
|
||||
* (at your option) any later version.
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
using System.IO;
|
||||
|
||||
namespace Server
|
||||
{
|
||||
public sealed class BinaryMemoryWriter : BinaryFileWriter
|
||||
{
|
||||
private static byte[] indexBuffer;
|
||||
private readonly MemoryStream stream;
|
||||
|
||||
public BinaryMemoryWriter()
|
||||
: base(new MemoryStream(512), true) =>
|
||||
stream = UnderlyingStream as MemoryStream;
|
||||
|
||||
protected override int BufferSize => 512;
|
||||
|
||||
public int CommitTo(
|
||||
SequentialFileWriterStream dataFile, SequentialFileWriterStream indexFile, int typeCode, uint serial
|
||||
)
|
||||
{
|
||||
Flush();
|
||||
|
||||
var buffer = stream.GetBuffer();
|
||||
var length = (int)stream.Length;
|
||||
|
||||
var position = dataFile.Position;
|
||||
|
||||
dataFile.Write(buffer, 0, length);
|
||||
|
||||
indexBuffer ??= new byte[20];
|
||||
|
||||
indexBuffer[0] = (byte)typeCode;
|
||||
indexBuffer[1] = (byte)(typeCode >> 8);
|
||||
indexBuffer[2] = (byte)(typeCode >> 16);
|
||||
indexBuffer[3] = (byte)(typeCode >> 24);
|
||||
|
||||
indexBuffer[4] = (byte)serial;
|
||||
indexBuffer[5] = (byte)(serial >> 8);
|
||||
indexBuffer[6] = (byte)(serial >> 16);
|
||||
indexBuffer[7] = (byte)(serial >> 24);
|
||||
|
||||
indexBuffer[8] = (byte)position;
|
||||
indexBuffer[9] = (byte)(position >> 8);
|
||||
indexBuffer[10] = (byte)(position >> 16);
|
||||
indexBuffer[11] = (byte)(position >> 24);
|
||||
indexBuffer[12] = (byte)(position >> 32);
|
||||
indexBuffer[13] = (byte)(position >> 40);
|
||||
indexBuffer[14] = (byte)(position >> 48);
|
||||
indexBuffer[15] = (byte)(position >> 56);
|
||||
|
||||
indexBuffer[16] = (byte)length;
|
||||
indexBuffer[17] = (byte)(length >> 8);
|
||||
indexBuffer[18] = (byte)(length >> 16);
|
||||
indexBuffer[19] = (byte)(length >> 24);
|
||||
|
||||
indexFile.Write(indexBuffer, 0, indexBuffer.Length);
|
||||
|
||||
stream.SetLength(0);
|
||||
|
||||
return length;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,47 +1,48 @@
|
|||
/***************************************************************************
|
||||
* DualSaveStrategy.cs
|
||||
* -------------------
|
||||
* begin : May 1, 2002
|
||||
* copyright : (C) The RunUO Software Team
|
||||
* email : info@runuo.com
|
||||
*
|
||||
* $Id$
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
/***************************************************************************
|
||||
*
|
||||
* This program is free software; you can redistribute it and/or modify
|
||||
* it under the terms of the GNU General Public License as published by
|
||||
* the Free Software Foundation; either version 2 of the License, or
|
||||
* (at your option) any later version.
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
using System.Threading;
|
||||
|
||||
namespace Server
|
||||
{
|
||||
public sealed class DualSaveStrategy : StandardSaveStrategy
|
||||
{
|
||||
public override string Name => "Dual";
|
||||
|
||||
public override void Save(bool permitBackgroundWrite)
|
||||
{
|
||||
PermitBackgroundWrite = permitBackgroundWrite;
|
||||
|
||||
var saveThread = new Thread(SaveItems);
|
||||
|
||||
saveThread.Name = "Item Save Subset";
|
||||
saveThread.Start();
|
||||
|
||||
SaveMobiles();
|
||||
SaveGuilds();
|
||||
|
||||
saveThread.Join();
|
||||
|
||||
if (permitBackgroundWrite && UseSequentialWriters) // If we're permitted to write in the background, but we don't anyways, then notify.
|
||||
World.NotifyDiskWriteComplete();
|
||||
}
|
||||
}
|
||||
}
|
||||
/***************************************************************************
|
||||
* DualSaveStrategy.cs
|
||||
* -------------------
|
||||
* begin : May 1, 2002
|
||||
* copyright : (C) The RunUO Software Team
|
||||
* email : info@runuo.com
|
||||
*
|
||||
* $Id$
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
/***************************************************************************
|
||||
*
|
||||
* This program is free software; you can redistribute it and/or modify
|
||||
* it under the terms of the GNU General Public License as published by
|
||||
* the Free Software Foundation; either version 2 of the License, or
|
||||
* (at your option) any later version.
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
using System.Threading;
|
||||
|
||||
namespace Server
|
||||
{
|
||||
public sealed class DualSaveStrategy : StandardSaveStrategy
|
||||
{
|
||||
public override string Name => "Dual";
|
||||
|
||||
public override void Save(bool permitBackgroundWrite)
|
||||
{
|
||||
PermitBackgroundWrite = permitBackgroundWrite;
|
||||
|
||||
var saveThread = new Thread(SaveItems);
|
||||
|
||||
saveThread.Name = "Item Save Subset";
|
||||
saveThread.Start();
|
||||
|
||||
SaveMobiles();
|
||||
SaveGuilds();
|
||||
|
||||
saveThread.Join();
|
||||
|
||||
if (permitBackgroundWrite && UseSequentialWriters
|
||||
) // If we're permitted to write in the background, but we don't anyways, then notify.
|
||||
World.NotifyDiskWriteComplete();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,281 +1,297 @@
|
|||
/***************************************************************************
|
||||
* DynamicSaveStrategy.cs
|
||||
* -------------------
|
||||
* begin : December 16, 2010
|
||||
* copyright : (C) The RunUO Software Team
|
||||
* email : info@runuo.com
|
||||
*
|
||||
* $Id$
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
/***************************************************************************
|
||||
*
|
||||
* This program is free software; you can redistribute it and/or modify
|
||||
* it under the terms of the GNU General Public License as published by
|
||||
* the Free Software Foundation; either version 2 of the License, or
|
||||
* (at your option) any later version.
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
using System;
|
||||
using System.Collections.Concurrent;
|
||||
using System.Collections.Generic;
|
||||
using System.Threading.Tasks;
|
||||
using Server.Guilds;
|
||||
|
||||
namespace Server
|
||||
{
|
||||
public sealed class DynamicSaveStrategy : SaveStrategy
|
||||
{
|
||||
private readonly ConcurrentBag<Item> _decayBag;
|
||||
private SequentialFileWriterStream _guildData, _guildIndex;
|
||||
private readonly BlockingCollection<QueuedMemoryWriter> _guildThreadWriters;
|
||||
|
||||
private SequentialFileWriterStream _itemData, _itemIndex;
|
||||
|
||||
private readonly BlockingCollection<QueuedMemoryWriter> _itemThreadWriters;
|
||||
|
||||
private SequentialFileWriterStream _mobileData, _mobileIndex;
|
||||
private readonly BlockingCollection<QueuedMemoryWriter> _mobileThreadWriters;
|
||||
|
||||
public DynamicSaveStrategy()
|
||||
{
|
||||
_decayBag = new ConcurrentBag<Item>();
|
||||
_itemThreadWriters = new BlockingCollection<QueuedMemoryWriter>();
|
||||
_mobileThreadWriters = new BlockingCollection<QueuedMemoryWriter>();
|
||||
_guildThreadWriters = new BlockingCollection<QueuedMemoryWriter>();
|
||||
}
|
||||
|
||||
public override string Name => "Dynamic";
|
||||
|
||||
public override void Save(bool permitBackgroundWrite)
|
||||
{
|
||||
OpenFiles();
|
||||
|
||||
var saveTasks = new Task[3];
|
||||
|
||||
saveTasks[0] = SaveItems();
|
||||
saveTasks[1] = SaveMobiles();
|
||||
saveTasks[2] = SaveGuilds();
|
||||
|
||||
SaveTypeDatabases();
|
||||
|
||||
if (permitBackgroundWrite)
|
||||
{
|
||||
#pragma warning restore CA2008 // Do not create tasks without passing a TaskScheduler *
|
||||
// This option makes it finish the writing to disk in the background, continuing even after Save() returns.
|
||||
Task.Factory.ContinueWhenAll(saveTasks, _ =>
|
||||
{
|
||||
CloseFiles();
|
||||
|
||||
World.NotifyDiskWriteComplete();
|
||||
});
|
||||
#pragma warning restore CA2008 // Do not create tasks without passing a TaskScheduler *
|
||||
}
|
||||
else
|
||||
{
|
||||
Task.WaitAll(saveTasks); // Waits for the completion of all of the tasks(committing to disk)
|
||||
CloseFiles();
|
||||
}
|
||||
}
|
||||
|
||||
private Task StartCommitTask(BlockingCollection<QueuedMemoryWriter> threadWriter, SequentialFileWriterStream data,
|
||||
SequentialFileWriterStream index)
|
||||
{
|
||||
#pragma warning disable CA2008 // Do not create tasks without passing a TaskScheduler *
|
||||
var commitTask = Task.Factory.StartNew(() =>
|
||||
{
|
||||
while (!threadWriter.IsCompleted)
|
||||
{
|
||||
QueuedMemoryWriter writer;
|
||||
|
||||
try
|
||||
{
|
||||
writer = threadWriter.Take();
|
||||
}
|
||||
catch (InvalidOperationException)
|
||||
{
|
||||
// Per MSDN, it's fine if we're here, successful completion of adding can rarely put us into this state.
|
||||
break;
|
||||
}
|
||||
|
||||
writer.CommitTo(data, index);
|
||||
}
|
||||
});
|
||||
#pragma warning restore CA2008 // Do not create tasks without passing a TaskScheduler *
|
||||
|
||||
return commitTask;
|
||||
}
|
||||
|
||||
private Task SaveItems()
|
||||
{
|
||||
// Start the blocking consumer; this runs in background.
|
||||
var commitTask = StartCommitTask(_itemThreadWriters, _itemData, _itemIndex);
|
||||
|
||||
IEnumerable<Item> items = World.Items.Values;
|
||||
|
||||
// Start the producer.
|
||||
Parallel.ForEach(items, () => new QueuedMemoryWriter(),
|
||||
(item, state, writer) =>
|
||||
{
|
||||
var startPosition = writer.Position;
|
||||
|
||||
item.Serialize(writer);
|
||||
|
||||
var size = (int)(writer.Position - startPosition);
|
||||
|
||||
writer.QueueForIndex(item, size);
|
||||
|
||||
if (item.Decays && item.Parent == null && item.Map != Map.Internal &&
|
||||
DateTime.UtcNow > item.LastMoved + item.DecayTime) _decayBag.Add(item);
|
||||
|
||||
return writer;
|
||||
},
|
||||
writer =>
|
||||
{
|
||||
writer.Flush();
|
||||
|
||||
_itemThreadWriters.Add(writer);
|
||||
});
|
||||
|
||||
_itemThreadWriters.CompleteAdding(); // We only get here after the Parallel.ForEach completes. Lets our task
|
||||
|
||||
return commitTask;
|
||||
}
|
||||
|
||||
private Task SaveMobiles()
|
||||
{
|
||||
// Start the blocking consumer; this runs in background.
|
||||
var commitTask = StartCommitTask(_mobileThreadWriters, _mobileData, _mobileIndex);
|
||||
|
||||
IEnumerable<Mobile> mobiles = World.Mobiles.Values;
|
||||
|
||||
// Start the producer.
|
||||
Parallel.ForEach(mobiles, () => new QueuedMemoryWriter(),
|
||||
(mobile, state, writer) =>
|
||||
{
|
||||
var startPosition = writer.Position;
|
||||
|
||||
mobile.Serialize(writer);
|
||||
|
||||
var size = (int)(writer.Position - startPosition);
|
||||
|
||||
writer.QueueForIndex(mobile, size);
|
||||
|
||||
return writer;
|
||||
},
|
||||
writer =>
|
||||
{
|
||||
writer.Flush();
|
||||
|
||||
_mobileThreadWriters.Add(writer);
|
||||
});
|
||||
|
||||
_mobileThreadWriters
|
||||
.CompleteAdding(); // We only get here after the Parallel.ForEach completes. Lets our task tell the consumer that we're done
|
||||
|
||||
return commitTask;
|
||||
}
|
||||
|
||||
private Task SaveGuilds()
|
||||
{
|
||||
// Start the blocking consumer; this runs in background.
|
||||
var commitTask = StartCommitTask(_guildThreadWriters, _guildData, _guildIndex);
|
||||
|
||||
IEnumerable<BaseGuild> guilds = BaseGuild.List.Values;
|
||||
|
||||
// Start the producer.
|
||||
Parallel.ForEach(guilds, () => new QueuedMemoryWriter(),
|
||||
(guild, state, writer) =>
|
||||
{
|
||||
var startPosition = writer.Position;
|
||||
|
||||
guild.Serialize(writer);
|
||||
|
||||
var size = (int)(writer.Position - startPosition);
|
||||
|
||||
writer.QueueForIndex(guild, size);
|
||||
|
||||
return writer;
|
||||
},
|
||||
writer =>
|
||||
{
|
||||
writer.Flush();
|
||||
|
||||
_guildThreadWriters.Add(writer);
|
||||
});
|
||||
|
||||
_guildThreadWriters.CompleteAdding(); // We only get here after the Parallel.ForEach completes. Lets our task
|
||||
|
||||
return commitTask;
|
||||
}
|
||||
|
||||
public override void ProcessDecay()
|
||||
{
|
||||
while (_decayBag.TryTake(out var item))
|
||||
if (item.OnDecay())
|
||||
item.Delete();
|
||||
}
|
||||
|
||||
private void OpenFiles()
|
||||
{
|
||||
_itemData = new SequentialFileWriterStream(World.ItemDataPath);
|
||||
_itemIndex = new SequentialFileWriterStream(World.ItemIndexPath);
|
||||
|
||||
_mobileData = new SequentialFileWriterStream(World.MobileDataPath);
|
||||
_mobileIndex = new SequentialFileWriterStream(World.MobileIndexPath);
|
||||
|
||||
_guildData = new SequentialFileWriterStream(World.GuildDataPath);
|
||||
_guildIndex = new SequentialFileWriterStream(World.GuildIndexPath);
|
||||
|
||||
WriteCount(_itemIndex, World.Items.Count);
|
||||
WriteCount(_mobileIndex, World.Mobiles.Count);
|
||||
WriteCount(_guildIndex, BaseGuild.List.Count);
|
||||
}
|
||||
|
||||
private void CloseFiles()
|
||||
{
|
||||
_itemData.Close();
|
||||
_itemIndex.Close();
|
||||
|
||||
_mobileData.Close();
|
||||
_mobileIndex.Close();
|
||||
|
||||
_guildData.Close();
|
||||
_guildIndex.Close();
|
||||
}
|
||||
|
||||
private void WriteCount(SequentialFileWriterStream indexFile, int count)
|
||||
{
|
||||
// Equiv to GenericWriter.Write( (int)count );
|
||||
var buffer = new byte[4];
|
||||
|
||||
buffer[0] = (byte)count;
|
||||
buffer[1] = (byte)(count >> 8);
|
||||
buffer[2] = (byte)(count >> 16);
|
||||
buffer[3] = (byte)(count >> 24);
|
||||
|
||||
indexFile.Write(buffer, 0, buffer.Length);
|
||||
}
|
||||
|
||||
private void SaveTypeDatabases()
|
||||
{
|
||||
SaveTypeDatabase(World.ItemTypesPath, World.m_ItemTypes);
|
||||
SaveTypeDatabase(World.MobileTypesPath, World.m_MobileTypes);
|
||||
}
|
||||
|
||||
private void SaveTypeDatabase(string path, List<Type> types)
|
||||
{
|
||||
var bfw = new BinaryFileWriter(path, false);
|
||||
|
||||
bfw.Write(types.Count);
|
||||
|
||||
foreach (var type in types) bfw.Write(type.FullName);
|
||||
|
||||
bfw.Flush();
|
||||
|
||||
bfw.Close();
|
||||
}
|
||||
}
|
||||
}
|
||||
/***************************************************************************
|
||||
* DynamicSaveStrategy.cs
|
||||
* -------------------
|
||||
* begin : December 16, 2010
|
||||
* copyright : (C) The RunUO Software Team
|
||||
* email : info@runuo.com
|
||||
*
|
||||
* $Id$
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
/***************************************************************************
|
||||
*
|
||||
* This program is free software; you can redistribute it and/or modify
|
||||
* it under the terms of the GNU General Public License as published by
|
||||
* the Free Software Foundation; either version 2 of the License, or
|
||||
* (at your option) any later version.
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
using System;
|
||||
using System.Collections.Concurrent;
|
||||
using System.Collections.Generic;
|
||||
using System.Threading.Tasks;
|
||||
using Server.Guilds;
|
||||
|
||||
namespace Server
|
||||
{
|
||||
public sealed class DynamicSaveStrategy : SaveStrategy
|
||||
{
|
||||
private readonly ConcurrentBag<Item> _decayBag;
|
||||
private SequentialFileWriterStream _guildData, _guildIndex;
|
||||
private readonly BlockingCollection<QueuedMemoryWriter> _guildThreadWriters;
|
||||
|
||||
private SequentialFileWriterStream _itemData, _itemIndex;
|
||||
|
||||
private readonly BlockingCollection<QueuedMemoryWriter> _itemThreadWriters;
|
||||
|
||||
private SequentialFileWriterStream _mobileData, _mobileIndex;
|
||||
private readonly BlockingCollection<QueuedMemoryWriter> _mobileThreadWriters;
|
||||
|
||||
public DynamicSaveStrategy()
|
||||
{
|
||||
_decayBag = new ConcurrentBag<Item>();
|
||||
_itemThreadWriters = new BlockingCollection<QueuedMemoryWriter>();
|
||||
_mobileThreadWriters = new BlockingCollection<QueuedMemoryWriter>();
|
||||
_guildThreadWriters = new BlockingCollection<QueuedMemoryWriter>();
|
||||
}
|
||||
|
||||
public override string Name => "Dynamic";
|
||||
|
||||
public override void Save(bool permitBackgroundWrite)
|
||||
{
|
||||
OpenFiles();
|
||||
|
||||
var saveTasks = new Task[3];
|
||||
|
||||
saveTasks[0] = SaveItems();
|
||||
saveTasks[1] = SaveMobiles();
|
||||
saveTasks[2] = SaveGuilds();
|
||||
|
||||
SaveTypeDatabases();
|
||||
|
||||
if (permitBackgroundWrite)
|
||||
{
|
||||
#pragma warning restore CA2008 // Do not create tasks without passing a TaskScheduler *
|
||||
// This option makes it finish the writing to disk in the background, continuing even after Save() returns.
|
||||
Task.Factory.ContinueWhenAll(
|
||||
saveTasks,
|
||||
_ =>
|
||||
{
|
||||
CloseFiles();
|
||||
|
||||
World.NotifyDiskWriteComplete();
|
||||
}
|
||||
);
|
||||
#pragma warning restore CA2008 // Do not create tasks without passing a TaskScheduler *
|
||||
}
|
||||
else
|
||||
{
|
||||
Task.WaitAll(saveTasks); // Waits for the completion of all of the tasks(committing to disk)
|
||||
CloseFiles();
|
||||
}
|
||||
}
|
||||
|
||||
private Task StartCommitTask(
|
||||
BlockingCollection<QueuedMemoryWriter> threadWriter, SequentialFileWriterStream data,
|
||||
SequentialFileWriterStream index
|
||||
)
|
||||
{
|
||||
#pragma warning disable CA2008 // Do not create tasks without passing a TaskScheduler *
|
||||
var commitTask = Task.Factory.StartNew(
|
||||
() =>
|
||||
{
|
||||
while (!threadWriter.IsCompleted)
|
||||
{
|
||||
QueuedMemoryWriter writer;
|
||||
|
||||
try
|
||||
{
|
||||
writer = threadWriter.Take();
|
||||
}
|
||||
catch (InvalidOperationException)
|
||||
{
|
||||
// Per MSDN, it's fine if we're here, successful completion of adding can rarely put us into this state.
|
||||
break;
|
||||
}
|
||||
|
||||
writer.CommitTo(data, index);
|
||||
}
|
||||
}
|
||||
);
|
||||
#pragma warning restore CA2008 // Do not create tasks without passing a TaskScheduler *
|
||||
|
||||
return commitTask;
|
||||
}
|
||||
|
||||
private Task SaveItems()
|
||||
{
|
||||
// Start the blocking consumer; this runs in background.
|
||||
var commitTask = StartCommitTask(_itemThreadWriters, _itemData, _itemIndex);
|
||||
|
||||
IEnumerable<Item> items = World.Items.Values;
|
||||
|
||||
// Start the producer.
|
||||
Parallel.ForEach(
|
||||
items,
|
||||
() => new QueuedMemoryWriter(),
|
||||
(item, state, writer) =>
|
||||
{
|
||||
var startPosition = writer.Position;
|
||||
|
||||
item.Serialize(writer);
|
||||
|
||||
var size = (int)(writer.Position - startPosition);
|
||||
|
||||
writer.QueueForIndex(item, size);
|
||||
|
||||
if (item.Decays && item.Parent == null && item.Map != Map.Internal &&
|
||||
DateTime.UtcNow > item.LastMoved + item.DecayTime) _decayBag.Add(item);
|
||||
|
||||
return writer;
|
||||
},
|
||||
writer =>
|
||||
{
|
||||
writer.Flush();
|
||||
|
||||
_itemThreadWriters.Add(writer);
|
||||
}
|
||||
);
|
||||
|
||||
_itemThreadWriters.CompleteAdding(); // We only get here after the Parallel.ForEach completes. Lets our task
|
||||
|
||||
return commitTask;
|
||||
}
|
||||
|
||||
private Task SaveMobiles()
|
||||
{
|
||||
// Start the blocking consumer; this runs in background.
|
||||
var commitTask = StartCommitTask(_mobileThreadWriters, _mobileData, _mobileIndex);
|
||||
|
||||
IEnumerable<Mobile> mobiles = World.Mobiles.Values;
|
||||
|
||||
// Start the producer.
|
||||
Parallel.ForEach(
|
||||
mobiles,
|
||||
() => new QueuedMemoryWriter(),
|
||||
(mobile, state, writer) =>
|
||||
{
|
||||
var startPosition = writer.Position;
|
||||
|
||||
mobile.Serialize(writer);
|
||||
|
||||
var size = (int)(writer.Position - startPosition);
|
||||
|
||||
writer.QueueForIndex(mobile, size);
|
||||
|
||||
return writer;
|
||||
},
|
||||
writer =>
|
||||
{
|
||||
writer.Flush();
|
||||
|
||||
_mobileThreadWriters.Add(writer);
|
||||
}
|
||||
);
|
||||
|
||||
_mobileThreadWriters
|
||||
.CompleteAdding(); // We only get here after the Parallel.ForEach completes. Lets our task tell the consumer that we're done
|
||||
|
||||
return commitTask;
|
||||
}
|
||||
|
||||
private Task SaveGuilds()
|
||||
{
|
||||
// Start the blocking consumer; this runs in background.
|
||||
var commitTask = StartCommitTask(_guildThreadWriters, _guildData, _guildIndex);
|
||||
|
||||
IEnumerable<BaseGuild> guilds = BaseGuild.List.Values;
|
||||
|
||||
// Start the producer.
|
||||
Parallel.ForEach(
|
||||
guilds,
|
||||
() => new QueuedMemoryWriter(),
|
||||
(guild, state, writer) =>
|
||||
{
|
||||
var startPosition = writer.Position;
|
||||
|
||||
guild.Serialize(writer);
|
||||
|
||||
var size = (int)(writer.Position - startPosition);
|
||||
|
||||
writer.QueueForIndex(guild, size);
|
||||
|
||||
return writer;
|
||||
},
|
||||
writer =>
|
||||
{
|
||||
writer.Flush();
|
||||
|
||||
_guildThreadWriters.Add(writer);
|
||||
}
|
||||
);
|
||||
|
||||
_guildThreadWriters.CompleteAdding(); // We only get here after the Parallel.ForEach completes. Lets our task
|
||||
|
||||
return commitTask;
|
||||
}
|
||||
|
||||
public override void ProcessDecay()
|
||||
{
|
||||
while (_decayBag.TryTake(out var item))
|
||||
if (item.OnDecay())
|
||||
item.Delete();
|
||||
}
|
||||
|
||||
private void OpenFiles()
|
||||
{
|
||||
_itemData = new SequentialFileWriterStream(World.ItemDataPath);
|
||||
_itemIndex = new SequentialFileWriterStream(World.ItemIndexPath);
|
||||
|
||||
_mobileData = new SequentialFileWriterStream(World.MobileDataPath);
|
||||
_mobileIndex = new SequentialFileWriterStream(World.MobileIndexPath);
|
||||
|
||||
_guildData = new SequentialFileWriterStream(World.GuildDataPath);
|
||||
_guildIndex = new SequentialFileWriterStream(World.GuildIndexPath);
|
||||
|
||||
WriteCount(_itemIndex, World.Items.Count);
|
||||
WriteCount(_mobileIndex, World.Mobiles.Count);
|
||||
WriteCount(_guildIndex, BaseGuild.List.Count);
|
||||
}
|
||||
|
||||
private void CloseFiles()
|
||||
{
|
||||
_itemData.Close();
|
||||
_itemIndex.Close();
|
||||
|
||||
_mobileData.Close();
|
||||
_mobileIndex.Close();
|
||||
|
||||
_guildData.Close();
|
||||
_guildIndex.Close();
|
||||
}
|
||||
|
||||
private void WriteCount(SequentialFileWriterStream indexFile, int count)
|
||||
{
|
||||
// Equiv to GenericWriter.Write( (int)count );
|
||||
var buffer = new byte[4];
|
||||
|
||||
buffer[0] = (byte)count;
|
||||
buffer[1] = (byte)(count >> 8);
|
||||
buffer[2] = (byte)(count >> 16);
|
||||
buffer[3] = (byte)(count >> 24);
|
||||
|
||||
indexFile.Write(buffer, 0, buffer.Length);
|
||||
}
|
||||
|
||||
private void SaveTypeDatabases()
|
||||
{
|
||||
SaveTypeDatabase(World.ItemTypesPath, World.m_ItemTypes);
|
||||
SaveTypeDatabase(World.MobileTypesPath, World.m_MobileTypes);
|
||||
}
|
||||
|
||||
private void SaveTypeDatabase(string path, List<Type> types)
|
||||
{
|
||||
var bfw = new BinaryFileWriter(path, false);
|
||||
|
||||
bfw.Write(types.Count);
|
||||
|
||||
foreach (var type in types) bfw.Write(type.FullName);
|
||||
|
||||
bfw.Flush();
|
||||
|
||||
bfw.Close();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,107 +1,108 @@
|
|||
/***************************************************************************
|
||||
* FileOperations.cs
|
||||
* -------------------
|
||||
* begin : May 1, 2002
|
||||
* copyright : (C) The RunUO Software Team
|
||||
* email : info@runuo.com
|
||||
*
|
||||
* $Id$
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
/***************************************************************************
|
||||
*
|
||||
* This program is free software; you can redistribute it and/or modify
|
||||
* it under the terms of the GNU General Public License as published by
|
||||
* the Free Software Foundation; either version 2 of the License, or
|
||||
* (at your option) any later version.
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
using System.IO;
|
||||
#if WINDOWS
|
||||
using System;
|
||||
using System.Runtime.InteropServices;
|
||||
using Microsoft.Win32.SafeHandles;
|
||||
#endif
|
||||
|
||||
namespace Server
|
||||
{
|
||||
public static class FileOperations
|
||||
{
|
||||
public const int KB = 1024;
|
||||
public const int MB = 1024 * KB;
|
||||
|
||||
public static int BufferSize { get; set; } = 1 * MB;
|
||||
|
||||
public static int Concurrency { get; set; } = 1;
|
||||
|
||||
#if WINDOWS
|
||||
public static bool Unbuffered { get; set; } = true;
|
||||
#endif
|
||||
|
||||
public static bool AreSynchronous => Concurrency < 1;
|
||||
|
||||
public static FileStream OpenSequentialStream(string path, FileMode mode, FileAccess access, FileShare share)
|
||||
{
|
||||
var options = FileOptions.SequentialScan;
|
||||
|
||||
if (Concurrency > 0)
|
||||
options |= FileOptions.Asynchronous;
|
||||
|
||||
#if !WINDOWS
|
||||
return new FileStream( path, mode, access, share, BufferSize, options );
|
||||
#else
|
||||
if (Unbuffered)
|
||||
options |= NoBuffering;
|
||||
else
|
||||
return new FileStream(path, mode, access, share, BufferSize, options);
|
||||
|
||||
var fileHandle =
|
||||
UnsafeNativeMethods.CreateFile(path, (int)access, share, IntPtr.Zero, mode, (int)options, IntPtr.Zero);
|
||||
|
||||
if (fileHandle.IsInvalid) throw new IOException();
|
||||
|
||||
return new UnbufferedFileStream(fileHandle, access, BufferSize, Concurrency > 0);
|
||||
#endif
|
||||
}
|
||||
|
||||
#if WINDOWS
|
||||
private class UnbufferedFileStream : FileStream
|
||||
{
|
||||
private readonly SafeFileHandle fileHandle;
|
||||
|
||||
public UnbufferedFileStream(SafeFileHandle fileHandle, FileAccess access, int bufferSize, bool isAsync)
|
||||
: base(fileHandle, access, bufferSize, isAsync) =>
|
||||
this.fileHandle = fileHandle;
|
||||
|
||||
public override void Write(byte[] array, int offset, int count)
|
||||
{
|
||||
base.Write(array, offset, BufferSize);
|
||||
}
|
||||
|
||||
public override IAsyncResult BeginWrite(byte[] array, int offset, int numBytes, AsyncCallback userCallback,
|
||||
object stateObject) =>
|
||||
base.BeginWrite(array, offset, BufferSize, userCallback, stateObject);
|
||||
|
||||
protected override void Dispose(bool disposing)
|
||||
{
|
||||
if (!fileHandle.IsClosed) fileHandle.Close();
|
||||
|
||||
base.Dispose(disposing);
|
||||
}
|
||||
}
|
||||
#endif
|
||||
|
||||
#if WINDOWS
|
||||
private const FileOptions NoBuffering = (FileOptions)0x20000000;
|
||||
|
||||
internal static class UnsafeNativeMethods
|
||||
{
|
||||
[DllImport("Kernel32", CharSet = CharSet.Unicode, SetLastError = true)]
|
||||
internal static extern SafeFileHandle CreateFile(string lpFileName, int dwDesiredAccess, FileShare dwShareMode,
|
||||
IntPtr securityAttrs, FileMode dwCreationDisposition, int dwFlagsAndAttributes, IntPtr hTemplateFile);
|
||||
}
|
||||
#endif
|
||||
}
|
||||
}
|
||||
/***************************************************************************
|
||||
* FileOperations.cs
|
||||
* -------------------
|
||||
* begin : May 1, 2002
|
||||
* copyright : (C) The RunUO Software Team
|
||||
* email : info@runuo.com
|
||||
*
|
||||
* $Id$
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
/***************************************************************************
|
||||
*
|
||||
* This program is free software; you can redistribute it and/or modify
|
||||
* it under the terms of the GNU General Public License as published by
|
||||
* the Free Software Foundation; either version 2 of the License, or
|
||||
* (at your option) any later version.
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
using System.IO;
|
||||
|
||||
#if WINDOWS
|
||||
using System;
|
||||
using System.Runtime.InteropServices;
|
||||
using Microsoft.Win32.SafeHandles;
|
||||
#endif
|
||||
|
||||
namespace Server
|
||||
{
|
||||
public static class FileOperations
|
||||
{
|
||||
public const int KB = 1024;
|
||||
public const int MB = 1024 * KB;
|
||||
|
||||
public static int BufferSize { get; set; } = 1 * MB;
|
||||
|
||||
public static int Concurrency { get; set; } = 1;
|
||||
|
||||
#if WINDOWS
|
||||
public static bool Unbuffered { get; set; } = true;
|
||||
#endif
|
||||
|
||||
public static bool AreSynchronous => Concurrency < 1;
|
||||
|
||||
public static FileStream OpenSequentialStream(string path, FileMode mode, FileAccess access, FileShare share)
|
||||
{
|
||||
var options = FileOptions.SequentialScan;
|
||||
|
||||
if (Concurrency > 0)
|
||||
options |= FileOptions.Asynchronous;
|
||||
|
||||
#if !WINDOWS
|
||||
return new FileStream(path, mode, access, share, BufferSize, options);
|
||||
#else
|
||||
if (Unbuffered)
|
||||
options |= NoBuffering;
|
||||
else
|
||||
return new FileStream(path, mode, access, share, BufferSize, options);
|
||||
|
||||
var fileHandle =
|
||||
UnsafeNativeMethods.CreateFile(path, (int)access, share, IntPtr.Zero, mode, (int)options, IntPtr.Zero);
|
||||
|
||||
if (fileHandle.IsInvalid) throw new IOException();
|
||||
|
||||
return new UnbufferedFileStream(fileHandle, access, BufferSize, Concurrency > 0);
|
||||
#endif
|
||||
}
|
||||
|
||||
#if WINDOWS
|
||||
private class UnbufferedFileStream : FileStream
|
||||
{
|
||||
private readonly SafeFileHandle fileHandle;
|
||||
|
||||
public UnbufferedFileStream(SafeFileHandle fileHandle, FileAccess access, int bufferSize, bool isAsync)
|
||||
: base(fileHandle, access, bufferSize, isAsync) =>
|
||||
this.fileHandle = fileHandle;
|
||||
|
||||
public override void Write(byte[] array, int offset, int count)
|
||||
{
|
||||
base.Write(array, offset, BufferSize);
|
||||
}
|
||||
|
||||
public override IAsyncResult BeginWrite(byte[] array, int offset, int numBytes, AsyncCallback userCallback,
|
||||
object stateObject) =>
|
||||
base.BeginWrite(array, offset, BufferSize, userCallback, stateObject);
|
||||
|
||||
protected override void Dispose(bool disposing)
|
||||
{
|
||||
if (!fileHandle.IsClosed) fileHandle.Close();
|
||||
|
||||
base.Dispose(disposing);
|
||||
}
|
||||
}
|
||||
#endif
|
||||
|
||||
#if WINDOWS
|
||||
private const FileOptions NoBuffering = (FileOptions)0x20000000;
|
||||
|
||||
internal static class UnsafeNativeMethods
|
||||
{
|
||||
[DllImport("Kernel32", CharSet = CharSet.Unicode, SetLastError = true)]
|
||||
internal static extern SafeFileHandle CreateFile(string lpFileName, int dwDesiredAccess, FileShare dwShareMode,
|
||||
IntPtr securityAttrs, FileMode dwCreationDisposition, int dwFlagsAndAttributes, IntPtr hTemplateFile);
|
||||
}
|
||||
#endif
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,228 +1,228 @@
|
|||
/***************************************************************************
|
||||
* FileQueue.cs
|
||||
* -------------------
|
||||
* begin : May 1, 2002
|
||||
* copyright : (C) The RunUO Software Team
|
||||
* email : info@runuo.com
|
||||
*
|
||||
* $Id$
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
/***************************************************************************
|
||||
*
|
||||
* This program is free software; you can redistribute it and/or modify
|
||||
* it under the terms of the GNU General Public License as published by
|
||||
* the Free Software Foundation; either version 2 of the License, or
|
||||
* (at your option) any later version.
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
using System;
|
||||
using System.Buffers;
|
||||
using System.Collections.Generic;
|
||||
using System.Threading;
|
||||
|
||||
namespace Server
|
||||
{
|
||||
public delegate void FileCommitCallback(FileQueue.Chunk chunk);
|
||||
|
||||
public sealed class FileQueue : IDisposable
|
||||
{
|
||||
private static readonly int bufferSize;
|
||||
|
||||
private readonly Chunk[] active;
|
||||
private int activeCount;
|
||||
private Page buffered;
|
||||
|
||||
private readonly FileCommitCallback callback;
|
||||
|
||||
private ManualResetEvent idle;
|
||||
|
||||
private readonly Queue<Page> pending;
|
||||
|
||||
private readonly object syncRoot;
|
||||
|
||||
static FileQueue() => bufferSize = FileOperations.BufferSize;
|
||||
|
||||
public FileQueue(int concurrentWrites, FileCommitCallback callback)
|
||||
{
|
||||
if (concurrentWrites < 1) throw new ArgumentOutOfRangeException(nameof(concurrentWrites));
|
||||
|
||||
if (bufferSize < 1)
|
||||
#pragma warning disable CA2208 // Instantiate argument exceptions correctly
|
||||
throw new ArgumentOutOfRangeException(nameof(FileOperations.BufferSize));
|
||||
#pragma warning restore CA2208 // Instantiate argument exceptions correctly
|
||||
|
||||
syncRoot = new object();
|
||||
|
||||
active = new Chunk[concurrentWrites];
|
||||
pending = new Queue<Page>();
|
||||
|
||||
this.callback = callback;
|
||||
|
||||
idle = new ManualResetEvent(true);
|
||||
}
|
||||
|
||||
public long Position { get; private set; }
|
||||
|
||||
public void Dispose()
|
||||
{
|
||||
if (idle != null)
|
||||
{
|
||||
idle.Close();
|
||||
idle = null;
|
||||
}
|
||||
}
|
||||
|
||||
private void Append(Page page)
|
||||
{
|
||||
lock (syncRoot)
|
||||
{
|
||||
if (activeCount == 0) idle.Reset();
|
||||
|
||||
++activeCount;
|
||||
|
||||
for (var slot = 0; slot < active.Length; ++slot)
|
||||
if (active[slot] == null)
|
||||
{
|
||||
active[slot] = new Chunk(this, slot, page.buffer, 0, page.length);
|
||||
|
||||
callback(active[slot]);
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
pending.Enqueue(page);
|
||||
}
|
||||
}
|
||||
|
||||
public void Flush()
|
||||
{
|
||||
if (buffered.buffer != null)
|
||||
{
|
||||
Append(buffered);
|
||||
|
||||
buffered.buffer = null;
|
||||
buffered.length = 0;
|
||||
}
|
||||
|
||||
/*lock ( syncRoot ) {
|
||||
if (pending.Count > 0 ) {
|
||||
idle.Reset();
|
||||
}
|
||||
|
||||
for ( int slot = 0; slot < active.Length && pending.Count > 0; ++slot ) {
|
||||
if (active[slot] == null ) {
|
||||
Page page = pending.Dequeue();
|
||||
|
||||
active[slot] = new Chunk( this, slot, page.buffer, 0, page.length );
|
||||
|
||||
++activeCount;
|
||||
|
||||
callback( active[slot] );
|
||||
}
|
||||
}
|
||||
}*/
|
||||
|
||||
idle.WaitOne();
|
||||
}
|
||||
|
||||
private void Commit(Chunk chunk, int slot)
|
||||
{
|
||||
if (slot < 0 || slot >= active.Length) throw new ArgumentOutOfRangeException(nameof(slot));
|
||||
|
||||
lock (syncRoot)
|
||||
{
|
||||
if (active[slot] != chunk) throw new ArgumentException("active slot is not the current chunk");
|
||||
|
||||
ArrayPool<byte>.Shared.Return(chunk.Buffer);
|
||||
|
||||
if (pending.Count > 0)
|
||||
{
|
||||
var page = pending.Dequeue();
|
||||
|
||||
active[slot] = new Chunk(this, slot, page.buffer, 0, page.length);
|
||||
|
||||
callback(active[slot]);
|
||||
}
|
||||
else
|
||||
{
|
||||
active[slot] = null;
|
||||
}
|
||||
|
||||
--activeCount;
|
||||
|
||||
if (activeCount == 0) idle.Set();
|
||||
}
|
||||
}
|
||||
|
||||
public void Enqueue(byte[] buffer, int offset, int size)
|
||||
{
|
||||
if (buffer == null) throw new ArgumentNullException(nameof(buffer));
|
||||
|
||||
if (offset < 0) throw new ArgumentOutOfRangeException(nameof(offset));
|
||||
if (size < 0) throw new ArgumentOutOfRangeException(nameof(size));
|
||||
if (buffer.Length - offset < size) throw new ArgumentOutOfRangeException(nameof(offset));
|
||||
|
||||
Position += size;
|
||||
|
||||
while (size > 0)
|
||||
{
|
||||
buffered.buffer ??= ArrayPool<byte>.Shared.Rent(bufferSize);
|
||||
|
||||
var page = buffered.buffer; // buffer page
|
||||
var pageSpace = page.Length - buffered.length; // available bytes in page
|
||||
var byteCount = size > pageSpace ? pageSpace : size; // how many bytes we can copy over
|
||||
|
||||
Buffer.BlockCopy(buffer, offset, page, buffered.length, byteCount);
|
||||
|
||||
buffered.length += byteCount;
|
||||
offset += byteCount;
|
||||
size -= byteCount;
|
||||
|
||||
if (buffered.length == page.Length)
|
||||
{
|
||||
// page full
|
||||
Append(buffered);
|
||||
|
||||
buffered.buffer = null;
|
||||
buffered.length = 0;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
public sealed class Chunk
|
||||
{
|
||||
private readonly FileQueue m_Owner;
|
||||
private readonly int m_Slot;
|
||||
|
||||
public Chunk(FileQueue owner, int slot, byte[] buffer, int offset, int size)
|
||||
{
|
||||
m_Owner = owner;
|
||||
m_Slot = slot;
|
||||
|
||||
Buffer = buffer;
|
||||
Offset = offset;
|
||||
Size = size;
|
||||
}
|
||||
|
||||
public byte[] Buffer { get; }
|
||||
|
||||
public int Offset { get; }
|
||||
|
||||
public int Size { get; }
|
||||
|
||||
public void Commit()
|
||||
{
|
||||
m_Owner.Commit(this, m_Slot);
|
||||
}
|
||||
}
|
||||
|
||||
private struct Page
|
||||
{
|
||||
public byte[] buffer;
|
||||
public int length;
|
||||
}
|
||||
}
|
||||
}
|
||||
/***************************************************************************
|
||||
* FileQueue.cs
|
||||
* -------------------
|
||||
* begin : May 1, 2002
|
||||
* copyright : (C) The RunUO Software Team
|
||||
* email : info@runuo.com
|
||||
*
|
||||
* $Id$
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
/***************************************************************************
|
||||
*
|
||||
* This program is free software; you can redistribute it and/or modify
|
||||
* it under the terms of the GNU General Public License as published by
|
||||
* the Free Software Foundation; either version 2 of the License, or
|
||||
* (at your option) any later version.
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
using System;
|
||||
using System.Buffers;
|
||||
using System.Collections.Generic;
|
||||
using System.Threading;
|
||||
|
||||
namespace Server
|
||||
{
|
||||
public delegate void FileCommitCallback(FileQueue.Chunk chunk);
|
||||
|
||||
public sealed class FileQueue : IDisposable
|
||||
{
|
||||
private static readonly int bufferSize;
|
||||
|
||||
private readonly Chunk[] active;
|
||||
|
||||
private readonly FileCommitCallback callback;
|
||||
|
||||
private readonly Queue<Page> pending;
|
||||
|
||||
private readonly object syncRoot;
|
||||
private int activeCount;
|
||||
private Page buffered;
|
||||
|
||||
private ManualResetEvent idle;
|
||||
|
||||
static FileQueue() => bufferSize = FileOperations.BufferSize;
|
||||
|
||||
public FileQueue(int concurrentWrites, FileCommitCallback callback)
|
||||
{
|
||||
if (concurrentWrites < 1) throw new ArgumentOutOfRangeException(nameof(concurrentWrites));
|
||||
|
||||
if (bufferSize < 1)
|
||||
#pragma warning disable CA2208 // Instantiate argument exceptions correctly
|
||||
throw new ArgumentOutOfRangeException(nameof(FileOperations.BufferSize));
|
||||
#pragma warning restore CA2208 // Instantiate argument exceptions correctly
|
||||
|
||||
syncRoot = new object();
|
||||
|
||||
active = new Chunk[concurrentWrites];
|
||||
pending = new Queue<Page>();
|
||||
|
||||
this.callback = callback;
|
||||
|
||||
idle = new ManualResetEvent(true);
|
||||
}
|
||||
|
||||
public long Position { get; private set; }
|
||||
|
||||
public void Dispose()
|
||||
{
|
||||
if (idle != null)
|
||||
{
|
||||
idle.Close();
|
||||
idle = null;
|
||||
}
|
||||
}
|
||||
|
||||
private void Append(Page page)
|
||||
{
|
||||
lock (syncRoot)
|
||||
{
|
||||
if (activeCount == 0) idle.Reset();
|
||||
|
||||
++activeCount;
|
||||
|
||||
for (var slot = 0; slot < active.Length; ++slot)
|
||||
if (active[slot] == null)
|
||||
{
|
||||
active[slot] = new Chunk(this, slot, page.buffer, 0, page.length);
|
||||
|
||||
callback(active[slot]);
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
pending.Enqueue(page);
|
||||
}
|
||||
}
|
||||
|
||||
public void Flush()
|
||||
{
|
||||
if (buffered.buffer != null)
|
||||
{
|
||||
Append(buffered);
|
||||
|
||||
buffered.buffer = null;
|
||||
buffered.length = 0;
|
||||
}
|
||||
|
||||
/*lock ( syncRoot ) {
|
||||
if (pending.Count > 0 ) {
|
||||
idle.Reset();
|
||||
}
|
||||
|
||||
for ( int slot = 0; slot < active.Length && pending.Count > 0; ++slot ) {
|
||||
if (active[slot] == null ) {
|
||||
Page page = pending.Dequeue();
|
||||
|
||||
active[slot] = new Chunk( this, slot, page.buffer, 0, page.length );
|
||||
|
||||
++activeCount;
|
||||
|
||||
callback( active[slot] );
|
||||
}
|
||||
}
|
||||
}*/
|
||||
|
||||
idle.WaitOne();
|
||||
}
|
||||
|
||||
private void Commit(Chunk chunk, int slot)
|
||||
{
|
||||
if (slot < 0 || slot >= active.Length) throw new ArgumentOutOfRangeException(nameof(slot));
|
||||
|
||||
lock (syncRoot)
|
||||
{
|
||||
if (active[slot] != chunk) throw new ArgumentException("active slot is not the current chunk");
|
||||
|
||||
ArrayPool<byte>.Shared.Return(chunk.Buffer);
|
||||
|
||||
if (pending.Count > 0)
|
||||
{
|
||||
var page = pending.Dequeue();
|
||||
|
||||
active[slot] = new Chunk(this, slot, page.buffer, 0, page.length);
|
||||
|
||||
callback(active[slot]);
|
||||
}
|
||||
else
|
||||
{
|
||||
active[slot] = null;
|
||||
}
|
||||
|
||||
--activeCount;
|
||||
|
||||
if (activeCount == 0) idle.Set();
|
||||
}
|
||||
}
|
||||
|
||||
public void Enqueue(byte[] buffer, int offset, int size)
|
||||
{
|
||||
if (buffer == null) throw new ArgumentNullException(nameof(buffer));
|
||||
|
||||
if (offset < 0) throw new ArgumentOutOfRangeException(nameof(offset));
|
||||
if (size < 0) throw new ArgumentOutOfRangeException(nameof(size));
|
||||
if (buffer.Length - offset < size) throw new ArgumentOutOfRangeException(nameof(offset));
|
||||
|
||||
Position += size;
|
||||
|
||||
while (size > 0)
|
||||
{
|
||||
buffered.buffer ??= ArrayPool<byte>.Shared.Rent(bufferSize);
|
||||
|
||||
var page = buffered.buffer; // buffer page
|
||||
var pageSpace = page.Length - buffered.length; // available bytes in page
|
||||
var byteCount = size > pageSpace ? pageSpace : size; // how many bytes we can copy over
|
||||
|
||||
Buffer.BlockCopy(buffer, offset, page, buffered.length, byteCount);
|
||||
|
||||
buffered.length += byteCount;
|
||||
offset += byteCount;
|
||||
size -= byteCount;
|
||||
|
||||
if (buffered.length == page.Length)
|
||||
{
|
||||
// page full
|
||||
Append(buffered);
|
||||
|
||||
buffered.buffer = null;
|
||||
buffered.length = 0;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
public sealed class Chunk
|
||||
{
|
||||
private readonly FileQueue m_Owner;
|
||||
private readonly int m_Slot;
|
||||
|
||||
public Chunk(FileQueue owner, int slot, byte[] buffer, int offset, int size)
|
||||
{
|
||||
m_Owner = owner;
|
||||
m_Slot = slot;
|
||||
|
||||
Buffer = buffer;
|
||||
Offset = offset;
|
||||
Size = size;
|
||||
}
|
||||
|
||||
public byte[] Buffer { get; }
|
||||
|
||||
public int Offset { get; }
|
||||
|
||||
public int Size { get; }
|
||||
|
||||
public void Commit()
|
||||
{
|
||||
m_Owner.Commit(this, m_Slot);
|
||||
}
|
||||
}
|
||||
|
||||
private struct Page
|
||||
{
|
||||
public byte[] buffer;
|
||||
public int length;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,352 +1,354 @@
|
|||
/***************************************************************************
|
||||
* ParallelSaveStrategy.cs
|
||||
* -------------------
|
||||
* begin : May 1, 2002
|
||||
* copyright : (C) The RunUO Software Team
|
||||
* email : info@runuo.com
|
||||
*
|
||||
* $Id$
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
/***************************************************************************
|
||||
*
|
||||
* This program is free software; you can redistribute it and/or modify
|
||||
* it under the terms of the GNU General Public License as published by
|
||||
* the Free Software Foundation; either version 2 of the License, or
|
||||
* (at your option) any later version.
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
using System;
|
||||
using System.Collections;
|
||||
using System.Collections.Generic;
|
||||
using System.Threading;
|
||||
using Server.Guilds;
|
||||
|
||||
namespace Server
|
||||
{
|
||||
public sealed class ParallelSaveStrategy : SaveStrategy, IDisposable
|
||||
{
|
||||
private readonly Queue<Item> _decayQueue;
|
||||
|
||||
private Consumer[] consumers;
|
||||
private int cycle;
|
||||
|
||||
private bool finished;
|
||||
private SequentialFileWriterStream guildData, guildIndex;
|
||||
|
||||
private SequentialFileWriterStream itemData, itemIndex;
|
||||
|
||||
private SequentialFileWriterStream mobileData, mobileIndex;
|
||||
|
||||
private readonly int processorCount;
|
||||
|
||||
public ParallelSaveStrategy(int processorCount)
|
||||
{
|
||||
this.processorCount = processorCount;
|
||||
|
||||
_decayQueue = new Queue<Item>();
|
||||
}
|
||||
|
||||
public override string Name => "Parallel";
|
||||
|
||||
private int GetThreadCount() => processorCount - 1;
|
||||
|
||||
public override void Save(bool permitBackgroundWrite)
|
||||
{
|
||||
OpenFiles();
|
||||
|
||||
consumers = new Consumer[GetThreadCount()];
|
||||
|
||||
for (var i = 0; i < consumers.Length; ++i) consumers[i] = new Consumer(this, 256);
|
||||
|
||||
IEnumerable<ISerializable> collection = new Producer();
|
||||
|
||||
foreach (var value in collection)
|
||||
while (!Enqueue(value))
|
||||
if (!Commit())
|
||||
Thread.Sleep(0);
|
||||
|
||||
finished = true;
|
||||
|
||||
SaveTypeDatabases();
|
||||
|
||||
WaitHandle.WaitAll(
|
||||
Array.ConvertAll<Consumer, WaitHandle>(
|
||||
consumers,
|
||||
input => input.completionEvent));
|
||||
|
||||
Commit();
|
||||
|
||||
CloseFiles();
|
||||
}
|
||||
|
||||
public override void ProcessDecay()
|
||||
{
|
||||
while (_decayQueue.Count > 0)
|
||||
{
|
||||
var item = _decayQueue.Dequeue();
|
||||
|
||||
if (item.OnDecay()) item.Delete();
|
||||
}
|
||||
}
|
||||
|
||||
private void SaveTypeDatabases()
|
||||
{
|
||||
SaveTypeDatabase(World.ItemTypesPath, World.m_ItemTypes);
|
||||
SaveTypeDatabase(World.MobileTypesPath, World.m_MobileTypes);
|
||||
}
|
||||
|
||||
private void SaveTypeDatabase(string path, List<Type> types)
|
||||
{
|
||||
var bfw = new BinaryFileWriter(path, false);
|
||||
|
||||
bfw.Write(types.Count);
|
||||
|
||||
foreach (var type in types) bfw.Write(type.FullName);
|
||||
|
||||
bfw.Flush();
|
||||
|
||||
bfw.Close();
|
||||
}
|
||||
|
||||
private void OpenFiles()
|
||||
{
|
||||
itemData = new SequentialFileWriterStream(World.ItemDataPath);
|
||||
itemIndex = new SequentialFileWriterStream(World.ItemIndexPath);
|
||||
|
||||
mobileData = new SequentialFileWriterStream(World.MobileDataPath);
|
||||
mobileIndex = new SequentialFileWriterStream(World.MobileIndexPath);
|
||||
|
||||
guildData = new SequentialFileWriterStream(World.GuildDataPath);
|
||||
guildIndex = new SequentialFileWriterStream(World.GuildIndexPath);
|
||||
|
||||
WriteCount(itemIndex, World.Items.Count);
|
||||
WriteCount(mobileIndex, World.Mobiles.Count);
|
||||
WriteCount(guildIndex, BaseGuild.List.Count);
|
||||
}
|
||||
|
||||
private void WriteCount(SequentialFileWriterStream indexFile, int count)
|
||||
{
|
||||
var buffer = new byte[4];
|
||||
|
||||
buffer[0] = (byte)count;
|
||||
buffer[1] = (byte)(count >> 8);
|
||||
buffer[2] = (byte)(count >> 16);
|
||||
buffer[3] = (byte)(count >> 24);
|
||||
|
||||
indexFile.Write(buffer, 0, buffer.Length);
|
||||
}
|
||||
|
||||
private void CloseFiles()
|
||||
{
|
||||
itemData.Close();
|
||||
itemIndex.Close();
|
||||
|
||||
mobileData.Close();
|
||||
mobileIndex.Close();
|
||||
|
||||
guildData.Close();
|
||||
guildIndex.Close();
|
||||
|
||||
World.NotifyDiskWriteComplete();
|
||||
}
|
||||
|
||||
private void OnSerialized(ConsumableEntry entry)
|
||||
{
|
||||
var value = entry.value;
|
||||
var writer = entry.writer;
|
||||
|
||||
if (value is Item item)
|
||||
Save(item, writer);
|
||||
else if (value is Mobile mob)
|
||||
Save(mob, writer);
|
||||
else if (value is BaseGuild guild)
|
||||
Save(guild, writer);
|
||||
}
|
||||
|
||||
private void Save(Item item, BinaryMemoryWriter writer)
|
||||
{
|
||||
writer.CommitTo(itemData, itemIndex, item.TypeRef, item.Serial);
|
||||
|
||||
if (item.Decays && item.Parent == null && item.Map != Map.Internal &&
|
||||
DateTime.UtcNow > item.LastMoved + item.DecayTime) _decayQueue.Enqueue(item);
|
||||
}
|
||||
|
||||
private void Save(Mobile mob, BinaryMemoryWriter writer)
|
||||
{
|
||||
writer.CommitTo(mobileData, mobileIndex, mob.TypeRef, mob.Serial);
|
||||
}
|
||||
|
||||
private void Save(BaseGuild guild, BinaryMemoryWriter writer)
|
||||
{
|
||||
writer.CommitTo(guildData, guildIndex, 0, guild.Serial);
|
||||
}
|
||||
|
||||
private bool Enqueue(ISerializable value)
|
||||
{
|
||||
for (var i = 0; i < consumers.Length; ++i)
|
||||
{
|
||||
var consumer = consumers[cycle++ % consumers.Length];
|
||||
|
||||
if (consumer.tail - consumer.head < consumer.buffer.Length)
|
||||
{
|
||||
consumer.buffer[consumer.tail % consumer.buffer.Length].value = value;
|
||||
consumer.tail++;
|
||||
|
||||
return true;
|
||||
}
|
||||
}
|
||||
|
||||
return false;
|
||||
}
|
||||
|
||||
private bool Commit()
|
||||
{
|
||||
var committed = false;
|
||||
|
||||
for (var i = 0; i < consumers.Length; ++i)
|
||||
{
|
||||
var consumer = consumers[i];
|
||||
|
||||
while (consumer.head < consumer.done)
|
||||
{
|
||||
OnSerialized(consumer.buffer[consumer.head % consumer.buffer.Length]);
|
||||
consumer.head++;
|
||||
|
||||
committed = true;
|
||||
}
|
||||
}
|
||||
|
||||
return committed;
|
||||
}
|
||||
|
||||
private sealed class Producer : IEnumerable<ISerializable>
|
||||
{
|
||||
private readonly IEnumerable<BaseGuild> guilds;
|
||||
private readonly IEnumerable<Item> items;
|
||||
private readonly IEnumerable<Mobile> mobiles;
|
||||
|
||||
public Producer()
|
||||
{
|
||||
items = World.Items.Values;
|
||||
mobiles = World.Mobiles.Values;
|
||||
guilds = BaseGuild.List.Values;
|
||||
}
|
||||
|
||||
public IEnumerator<ISerializable> GetEnumerator()
|
||||
{
|
||||
foreach (var item in items) yield return item;
|
||||
|
||||
foreach (var mob in mobiles) yield return mob;
|
||||
|
||||
foreach (var guild in guilds) yield return guild;
|
||||
}
|
||||
|
||||
IEnumerator IEnumerable.GetEnumerator() => throw new NotImplementedException();
|
||||
}
|
||||
|
||||
private struct ConsumableEntry
|
||||
{
|
||||
public ISerializable value;
|
||||
public BinaryMemoryWriter writer;
|
||||
}
|
||||
|
||||
private sealed class Consumer
|
||||
{
|
||||
public readonly ConsumableEntry[] buffer;
|
||||
|
||||
public readonly ManualResetEvent completionEvent;
|
||||
public int head, done, tail;
|
||||
private readonly ParallelSaveStrategy owner;
|
||||
|
||||
private readonly Thread thread;
|
||||
|
||||
public Consumer(ParallelSaveStrategy owner, int bufferSize)
|
||||
{
|
||||
this.owner = owner;
|
||||
|
||||
buffer = new ConsumableEntry[bufferSize];
|
||||
|
||||
for (var i = 0; i < buffer.Length; ++i) buffer[i].writer = new BinaryMemoryWriter();
|
||||
|
||||
completionEvent = new ManualResetEvent(false);
|
||||
|
||||
thread = new Thread(Processor);
|
||||
|
||||
thread.Name = "Parallel Serialization Thread";
|
||||
|
||||
thread.Start();
|
||||
}
|
||||
|
||||
private void Processor()
|
||||
{
|
||||
try
|
||||
{
|
||||
while (!owner.finished)
|
||||
{
|
||||
Process();
|
||||
Thread.Sleep(0);
|
||||
}
|
||||
|
||||
Process();
|
||||
|
||||
completionEvent.Set();
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
Console.WriteLine(ex);
|
||||
}
|
||||
}
|
||||
|
||||
private void Process()
|
||||
{
|
||||
ConsumableEntry entry;
|
||||
|
||||
while (done < tail)
|
||||
{
|
||||
entry = buffer[done % buffer.Length];
|
||||
|
||||
entry.value.Serialize(entry.writer);
|
||||
|
||||
++done;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private bool disposedValue = false; // To detect redundant calls
|
||||
|
||||
public void Dispose(bool disposing)
|
||||
{
|
||||
if (!disposedValue)
|
||||
{
|
||||
if (disposing)
|
||||
{
|
||||
// TODO: dispose managed state (managed objects).
|
||||
}
|
||||
|
||||
// TODO: free unmanaged resources (unmanaged objects) and override a finalizer below.
|
||||
// TODO: set large fields to null.
|
||||
|
||||
disposedValue = true;
|
||||
}
|
||||
}
|
||||
|
||||
// TODO: override a finalizer only if Dispose(bool disposing) above has code to free unmanaged resources.
|
||||
// ~ParallelSaveStrategy()
|
||||
// {
|
||||
// // Do not change this code. Put cleanup code in Dispose(bool disposing) above.
|
||||
// Dispose(false);
|
||||
// }
|
||||
|
||||
// This code added to correctly implement the disposable pattern.
|
||||
public void Dispose()
|
||||
{
|
||||
// Do not change this code. Put cleanup code in Dispose(bool disposing) above.
|
||||
Dispose(true);
|
||||
// TODO: uncomment the following line if the finalizer is overridden above.
|
||||
// GC.SuppressFinalize(this);
|
||||
}
|
||||
}
|
||||
}
|
||||
/***************************************************************************
|
||||
* ParallelSaveStrategy.cs
|
||||
* -------------------
|
||||
* begin : May 1, 2002
|
||||
* copyright : (C) The RunUO Software Team
|
||||
* email : info@runuo.com
|
||||
*
|
||||
* $Id$
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
/***************************************************************************
|
||||
*
|
||||
* This program is free software; you can redistribute it and/or modify
|
||||
* it under the terms of the GNU General Public License as published by
|
||||
* the Free Software Foundation; either version 2 of the License, or
|
||||
* (at your option) any later version.
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
using System;
|
||||
using System.Collections;
|
||||
using System.Collections.Generic;
|
||||
using System.Threading;
|
||||
using Server.Guilds;
|
||||
|
||||
namespace Server
|
||||
{
|
||||
public sealed class ParallelSaveStrategy : SaveStrategy, IDisposable
|
||||
{
|
||||
private readonly Queue<Item> _decayQueue;
|
||||
|
||||
private readonly int processorCount;
|
||||
|
||||
private Consumer[] consumers;
|
||||
private int cycle;
|
||||
|
||||
private bool disposedValue; // To detect redundant calls
|
||||
|
||||
private bool finished;
|
||||
private SequentialFileWriterStream guildData, guildIndex;
|
||||
|
||||
private SequentialFileWriterStream itemData, itemIndex;
|
||||
|
||||
private SequentialFileWriterStream mobileData, mobileIndex;
|
||||
|
||||
public ParallelSaveStrategy(int processorCount)
|
||||
{
|
||||
this.processorCount = processorCount;
|
||||
|
||||
_decayQueue = new Queue<Item>();
|
||||
}
|
||||
|
||||
public override string Name => "Parallel";
|
||||
|
||||
// TODO: override a finalizer only if Dispose(bool disposing) above has code to free unmanaged resources.
|
||||
// ~ParallelSaveStrategy()
|
||||
// {
|
||||
// // Do not change this code. Put cleanup code in Dispose(bool disposing) above.
|
||||
// Dispose(false);
|
||||
// }
|
||||
|
||||
// This code added to correctly implement the disposable pattern.
|
||||
public void Dispose()
|
||||
{
|
||||
// Do not change this code. Put cleanup code in Dispose(bool disposing) above.
|
||||
Dispose(true);
|
||||
// TODO: uncomment the following line if the finalizer is overridden above.
|
||||
// GC.SuppressFinalize(this);
|
||||
}
|
||||
|
||||
private int GetThreadCount() => processorCount - 1;
|
||||
|
||||
public override void Save(bool permitBackgroundWrite)
|
||||
{
|
||||
OpenFiles();
|
||||
|
||||
consumers = new Consumer[GetThreadCount()];
|
||||
|
||||
for (var i = 0; i < consumers.Length; ++i) consumers[i] = new Consumer(this, 256);
|
||||
|
||||
IEnumerable<ISerializable> collection = new Producer();
|
||||
|
||||
foreach (var value in collection)
|
||||
while (!Enqueue(value))
|
||||
if (!Commit())
|
||||
Thread.Sleep(0);
|
||||
|
||||
finished = true;
|
||||
|
||||
SaveTypeDatabases();
|
||||
|
||||
WaitHandle.WaitAll(
|
||||
Array.ConvertAll<Consumer, WaitHandle>(
|
||||
consumers,
|
||||
input => input.completionEvent
|
||||
)
|
||||
);
|
||||
|
||||
Commit();
|
||||
|
||||
CloseFiles();
|
||||
}
|
||||
|
||||
public override void ProcessDecay()
|
||||
{
|
||||
while (_decayQueue.Count > 0)
|
||||
{
|
||||
var item = _decayQueue.Dequeue();
|
||||
|
||||
if (item.OnDecay()) item.Delete();
|
||||
}
|
||||
}
|
||||
|
||||
private void SaveTypeDatabases()
|
||||
{
|
||||
SaveTypeDatabase(World.ItemTypesPath, World.m_ItemTypes);
|
||||
SaveTypeDatabase(World.MobileTypesPath, World.m_MobileTypes);
|
||||
}
|
||||
|
||||
private void SaveTypeDatabase(string path, List<Type> types)
|
||||
{
|
||||
var bfw = new BinaryFileWriter(path, false);
|
||||
|
||||
bfw.Write(types.Count);
|
||||
|
||||
foreach (var type in types) bfw.Write(type.FullName);
|
||||
|
||||
bfw.Flush();
|
||||
|
||||
bfw.Close();
|
||||
}
|
||||
|
||||
private void OpenFiles()
|
||||
{
|
||||
itemData = new SequentialFileWriterStream(World.ItemDataPath);
|
||||
itemIndex = new SequentialFileWriterStream(World.ItemIndexPath);
|
||||
|
||||
mobileData = new SequentialFileWriterStream(World.MobileDataPath);
|
||||
mobileIndex = new SequentialFileWriterStream(World.MobileIndexPath);
|
||||
|
||||
guildData = new SequentialFileWriterStream(World.GuildDataPath);
|
||||
guildIndex = new SequentialFileWriterStream(World.GuildIndexPath);
|
||||
|
||||
WriteCount(itemIndex, World.Items.Count);
|
||||
WriteCount(mobileIndex, World.Mobiles.Count);
|
||||
WriteCount(guildIndex, BaseGuild.List.Count);
|
||||
}
|
||||
|
||||
private void WriteCount(SequentialFileWriterStream indexFile, int count)
|
||||
{
|
||||
var buffer = new byte[4];
|
||||
|
||||
buffer[0] = (byte)count;
|
||||
buffer[1] = (byte)(count >> 8);
|
||||
buffer[2] = (byte)(count >> 16);
|
||||
buffer[3] = (byte)(count >> 24);
|
||||
|
||||
indexFile.Write(buffer, 0, buffer.Length);
|
||||
}
|
||||
|
||||
private void CloseFiles()
|
||||
{
|
||||
itemData.Close();
|
||||
itemIndex.Close();
|
||||
|
||||
mobileData.Close();
|
||||
mobileIndex.Close();
|
||||
|
||||
guildData.Close();
|
||||
guildIndex.Close();
|
||||
|
||||
World.NotifyDiskWriteComplete();
|
||||
}
|
||||
|
||||
private void OnSerialized(ConsumableEntry entry)
|
||||
{
|
||||
var value = entry.value;
|
||||
var writer = entry.writer;
|
||||
|
||||
if (value is Item item)
|
||||
Save(item, writer);
|
||||
else if (value is Mobile mob)
|
||||
Save(mob, writer);
|
||||
else if (value is BaseGuild guild)
|
||||
Save(guild, writer);
|
||||
}
|
||||
|
||||
private void Save(Item item, BinaryMemoryWriter writer)
|
||||
{
|
||||
writer.CommitTo(itemData, itemIndex, item.TypeRef, item.Serial);
|
||||
|
||||
if (item.Decays && item.Parent == null && item.Map != Map.Internal &&
|
||||
DateTime.UtcNow > item.LastMoved + item.DecayTime) _decayQueue.Enqueue(item);
|
||||
}
|
||||
|
||||
private void Save(Mobile mob, BinaryMemoryWriter writer)
|
||||
{
|
||||
writer.CommitTo(mobileData, mobileIndex, mob.TypeRef, mob.Serial);
|
||||
}
|
||||
|
||||
private void Save(BaseGuild guild, BinaryMemoryWriter writer)
|
||||
{
|
||||
writer.CommitTo(guildData, guildIndex, 0, guild.Serial);
|
||||
}
|
||||
|
||||
private bool Enqueue(ISerializable value)
|
||||
{
|
||||
for (var i = 0; i < consumers.Length; ++i)
|
||||
{
|
||||
var consumer = consumers[cycle++ % consumers.Length];
|
||||
|
||||
if (consumer.tail - consumer.head < consumer.buffer.Length)
|
||||
{
|
||||
consumer.buffer[consumer.tail % consumer.buffer.Length].value = value;
|
||||
consumer.tail++;
|
||||
|
||||
return true;
|
||||
}
|
||||
}
|
||||
|
||||
return false;
|
||||
}
|
||||
|
||||
private bool Commit()
|
||||
{
|
||||
var committed = false;
|
||||
|
||||
for (var i = 0; i < consumers.Length; ++i)
|
||||
{
|
||||
var consumer = consumers[i];
|
||||
|
||||
while (consumer.head < consumer.done)
|
||||
{
|
||||
OnSerialized(consumer.buffer[consumer.head % consumer.buffer.Length]);
|
||||
consumer.head++;
|
||||
|
||||
committed = true;
|
||||
}
|
||||
}
|
||||
|
||||
return committed;
|
||||
}
|
||||
|
||||
public void Dispose(bool disposing)
|
||||
{
|
||||
if (!disposedValue)
|
||||
{
|
||||
if (disposing)
|
||||
{
|
||||
// TODO: dispose managed state (managed objects).
|
||||
}
|
||||
|
||||
// TODO: free unmanaged resources (unmanaged objects) and override a finalizer below.
|
||||
// TODO: set large fields to null.
|
||||
|
||||
disposedValue = true;
|
||||
}
|
||||
}
|
||||
|
||||
private sealed class Producer : IEnumerable<ISerializable>
|
||||
{
|
||||
private readonly IEnumerable<BaseGuild> guilds;
|
||||
private readonly IEnumerable<Item> items;
|
||||
private readonly IEnumerable<Mobile> mobiles;
|
||||
|
||||
public Producer()
|
||||
{
|
||||
items = World.Items.Values;
|
||||
mobiles = World.Mobiles.Values;
|
||||
guilds = BaseGuild.List.Values;
|
||||
}
|
||||
|
||||
public IEnumerator<ISerializable> GetEnumerator()
|
||||
{
|
||||
foreach (var item in items) yield return item;
|
||||
|
||||
foreach (var mob in mobiles) yield return mob;
|
||||
|
||||
foreach (var guild in guilds) yield return guild;
|
||||
}
|
||||
|
||||
IEnumerator IEnumerable.GetEnumerator() => throw new NotImplementedException();
|
||||
}
|
||||
|
||||
private struct ConsumableEntry
|
||||
{
|
||||
public ISerializable value;
|
||||
public BinaryMemoryWriter writer;
|
||||
}
|
||||
|
||||
private sealed class Consumer
|
||||
{
|
||||
public readonly ConsumableEntry[] buffer;
|
||||
|
||||
public readonly ManualResetEvent completionEvent;
|
||||
private readonly ParallelSaveStrategy owner;
|
||||
|
||||
private readonly Thread thread;
|
||||
public int head, done, tail;
|
||||
|
||||
public Consumer(ParallelSaveStrategy owner, int bufferSize)
|
||||
{
|
||||
this.owner = owner;
|
||||
|
||||
buffer = new ConsumableEntry[bufferSize];
|
||||
|
||||
for (var i = 0; i < buffer.Length; ++i) buffer[i].writer = new BinaryMemoryWriter();
|
||||
|
||||
completionEvent = new ManualResetEvent(false);
|
||||
|
||||
thread = new Thread(Processor);
|
||||
|
||||
thread.Name = "Parallel Serialization Thread";
|
||||
|
||||
thread.Start();
|
||||
}
|
||||
|
||||
private void Processor()
|
||||
{
|
||||
try
|
||||
{
|
||||
while (!owner.finished)
|
||||
{
|
||||
Process();
|
||||
Thread.Sleep(0);
|
||||
}
|
||||
|
||||
Process();
|
||||
|
||||
completionEvent.Set();
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
Console.WriteLine(ex);
|
||||
}
|
||||
}
|
||||
|
||||
private void Process()
|
||||
{
|
||||
ConsumableEntry entry;
|
||||
|
||||
while (done < tail)
|
||||
{
|
||||
entry = buffer[done % buffer.Length];
|
||||
|
||||
entry.value.Serialize(entry.writer);
|
||||
|
||||
++done;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,99 +1,99 @@
|
|||
using System;
|
||||
using System.IO;
|
||||
|
||||
namespace Server
|
||||
{
|
||||
public static class Persistence
|
||||
{
|
||||
public static void Serialize(string path, Action<IGenericWriter> serializer)
|
||||
{
|
||||
Serialize(new FileInfo(path), serializer);
|
||||
}
|
||||
|
||||
public static void Serialize(FileInfo file, Action<IGenericWriter> serializer)
|
||||
{
|
||||
file.Refresh();
|
||||
|
||||
if (file.Directory?.Exists == false)
|
||||
file.Directory.Create();
|
||||
|
||||
if (!file.Exists) file.Create().Close();
|
||||
|
||||
file.Refresh();
|
||||
|
||||
using var fs = file.OpenWrite();
|
||||
var writer = new BinaryFileWriter(fs, true);
|
||||
|
||||
try
|
||||
{
|
||||
serializer(writer);
|
||||
}
|
||||
finally
|
||||
{
|
||||
writer.Flush();
|
||||
writer.Close();
|
||||
}
|
||||
}
|
||||
|
||||
public static void Deserialize(string path, Action<IGenericReader> deserializer)
|
||||
{
|
||||
Deserialize(path, deserializer, true);
|
||||
}
|
||||
|
||||
public static void Deserialize(FileInfo file, Action<IGenericReader> deserializer)
|
||||
{
|
||||
Deserialize(file, deserializer, true);
|
||||
}
|
||||
|
||||
public static void Deserialize(string path, Action<IGenericReader> deserializer, bool ensure)
|
||||
{
|
||||
Deserialize(new FileInfo(path), deserializer, ensure);
|
||||
}
|
||||
|
||||
public static void Deserialize(FileInfo file, Action<IGenericReader> deserializer, bool ensure)
|
||||
{
|
||||
file.Refresh();
|
||||
|
||||
if (file.Directory?.Exists == false)
|
||||
{
|
||||
if (!ensure)
|
||||
throw new DirectoryNotFoundException();
|
||||
|
||||
file.Directory.Create();
|
||||
}
|
||||
|
||||
if (!file.Exists)
|
||||
{
|
||||
if (!ensure)
|
||||
throw new FileNotFoundException
|
||||
{
|
||||
Source = file.FullName
|
||||
};
|
||||
|
||||
file.Create().Close();
|
||||
}
|
||||
|
||||
file.Refresh();
|
||||
|
||||
using var fs = file.OpenRead();
|
||||
var reader = new BinaryFileReader(new BinaryReader(fs));
|
||||
|
||||
try
|
||||
{
|
||||
deserializer(reader);
|
||||
}
|
||||
catch (EndOfStreamException eos)
|
||||
{
|
||||
if (file.Length > 0) Console.WriteLine("[Persistence]: {0}", eos);
|
||||
}
|
||||
catch (Exception e)
|
||||
{
|
||||
Console.WriteLine("[Persistence]: {0}", e);
|
||||
}
|
||||
finally
|
||||
{
|
||||
reader.Close();
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
using System;
|
||||
using System.IO;
|
||||
|
||||
namespace Server
|
||||
{
|
||||
public static class Persistence
|
||||
{
|
||||
public static void Serialize(string path, Action<IGenericWriter> serializer)
|
||||
{
|
||||
Serialize(new FileInfo(path), serializer);
|
||||
}
|
||||
|
||||
public static void Serialize(FileInfo file, Action<IGenericWriter> serializer)
|
||||
{
|
||||
file.Refresh();
|
||||
|
||||
if (file.Directory?.Exists == false)
|
||||
file.Directory.Create();
|
||||
|
||||
if (!file.Exists) file.Create().Close();
|
||||
|
||||
file.Refresh();
|
||||
|
||||
using var fs = file.OpenWrite();
|
||||
var writer = new BinaryFileWriter(fs, true);
|
||||
|
||||
try
|
||||
{
|
||||
serializer(writer);
|
||||
}
|
||||
finally
|
||||
{
|
||||
writer.Flush();
|
||||
writer.Close();
|
||||
}
|
||||
}
|
||||
|
||||
public static void Deserialize(string path, Action<IGenericReader> deserializer)
|
||||
{
|
||||
Deserialize(path, deserializer, true);
|
||||
}
|
||||
|
||||
public static void Deserialize(FileInfo file, Action<IGenericReader> deserializer)
|
||||
{
|
||||
Deserialize(file, deserializer, true);
|
||||
}
|
||||
|
||||
public static void Deserialize(string path, Action<IGenericReader> deserializer, bool ensure)
|
||||
{
|
||||
Deserialize(new FileInfo(path), deserializer, ensure);
|
||||
}
|
||||
|
||||
public static void Deserialize(FileInfo file, Action<IGenericReader> deserializer, bool ensure)
|
||||
{
|
||||
file.Refresh();
|
||||
|
||||
if (file.Directory?.Exists == false)
|
||||
{
|
||||
if (!ensure)
|
||||
throw new DirectoryNotFoundException();
|
||||
|
||||
file.Directory.Create();
|
||||
}
|
||||
|
||||
if (!file.Exists)
|
||||
{
|
||||
if (!ensure)
|
||||
throw new FileNotFoundException
|
||||
{
|
||||
Source = file.FullName
|
||||
};
|
||||
|
||||
file.Create().Close();
|
||||
}
|
||||
|
||||
file.Refresh();
|
||||
|
||||
using var fs = file.OpenRead();
|
||||
var reader = new BinaryFileReader(new BinaryReader(fs));
|
||||
|
||||
try
|
||||
{
|
||||
deserializer(reader);
|
||||
}
|
||||
catch (EndOfStreamException eos)
|
||||
{
|
||||
if (file.Length > 0) Console.WriteLine("[Persistence]: {0}", eos);
|
||||
}
|
||||
catch (Exception e)
|
||||
{
|
||||
Console.WriteLine("[Persistence]: {0}", e);
|
||||
}
|
||||
finally
|
||||
{
|
||||
reader.Close();
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,114 +1,114 @@
|
|||
/***************************************************************************
|
||||
* QueuedMemoryWriter.cs
|
||||
* -------------------
|
||||
* begin : December 16, 2010
|
||||
* copyright : (C) The RunUO Software Team
|
||||
* email : info@runuo.com
|
||||
*
|
||||
* $Id$
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
/***************************************************************************
|
||||
*
|
||||
* This program is free software; you can redistribute it and/or modify
|
||||
* it under the terms of the GNU General Public License as published by
|
||||
* the Free Software Foundation; either version 2 of the License, or
|
||||
* (at your option) any later version.
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
using System.Collections.Generic;
|
||||
using System.IO;
|
||||
|
||||
namespace Server
|
||||
{
|
||||
public sealed class QueuedMemoryWriter : BinaryFileWriter
|
||||
{
|
||||
private readonly MemoryStream _memStream;
|
||||
private readonly List<IndexInfo> _orderedIndexInfo = new List<IndexInfo>();
|
||||
|
||||
public QueuedMemoryWriter()
|
||||
: base(new MemoryStream(1024 * 1024), true) =>
|
||||
_memStream = UnderlyingStream as MemoryStream;
|
||||
|
||||
protected override int BufferSize => 512;
|
||||
|
||||
public void QueueForIndex(ISerializable serializable, int size)
|
||||
{
|
||||
IndexInfo info;
|
||||
|
||||
info.size = size;
|
||||
|
||||
info.typeCode = serializable.TypeRef; // For guilds, this will automagically be zero.
|
||||
info.serial = serializable.Serial;
|
||||
|
||||
_orderedIndexInfo.Add(info);
|
||||
}
|
||||
|
||||
public void CommitTo(SequentialFileWriterStream dataFile, SequentialFileWriterStream indexFile)
|
||||
{
|
||||
Flush();
|
||||
|
||||
var memLength = (int)_memStream.Position;
|
||||
|
||||
if (memLength > 0)
|
||||
{
|
||||
var memBuffer = _memStream.GetBuffer();
|
||||
|
||||
var actualPosition = dataFile.Position;
|
||||
|
||||
dataFile.Write(memBuffer, 0, memLength); // The buffer contains the data from many items.
|
||||
|
||||
// Console.WriteLine("Writing {0} bytes starting at {1}, with {2} things", memLength, actualPosition, _orderedIndexInfo.Count);
|
||||
|
||||
var indexBuffer = new byte[20];
|
||||
|
||||
// int indexWritten = _orderedIndexInfo.Count * indexBuffer.Length;
|
||||
// int totalWritten = memLength + indexWritten
|
||||
|
||||
for (var i = 0; i < _orderedIndexInfo.Count; i++)
|
||||
{
|
||||
var info = _orderedIndexInfo[i];
|
||||
|
||||
indexBuffer[0] = (byte)info.typeCode;
|
||||
indexBuffer[1] = (byte)(info.typeCode >> 8);
|
||||
indexBuffer[2] = (byte)(info.typeCode >> 16);
|
||||
indexBuffer[3] = (byte)(info.typeCode >> 24);
|
||||
|
||||
indexBuffer[4] = (byte)info.serial;
|
||||
indexBuffer[5] = (byte)(info.serial >> 8);
|
||||
indexBuffer[6] = (byte)(info.serial >> 16);
|
||||
indexBuffer[7] = (byte)(info.serial >> 24);
|
||||
|
||||
indexBuffer[8] = (byte)actualPosition;
|
||||
indexBuffer[9] = (byte)(actualPosition >> 8);
|
||||
indexBuffer[10] = (byte)(actualPosition >> 16);
|
||||
indexBuffer[11] = (byte)(actualPosition >> 24);
|
||||
indexBuffer[12] = (byte)(actualPosition >> 32);
|
||||
indexBuffer[13] = (byte)(actualPosition >> 40);
|
||||
indexBuffer[14] = (byte)(actualPosition >> 48);
|
||||
indexBuffer[15] = (byte)(actualPosition >> 56);
|
||||
|
||||
indexBuffer[16] = (byte)info.size;
|
||||
indexBuffer[17] = (byte)(info.size >> 8);
|
||||
indexBuffer[18] = (byte)(info.size >> 16);
|
||||
indexBuffer[19] = (byte)(info.size >> 24);
|
||||
|
||||
indexFile.Write(indexBuffer, 0, indexBuffer.Length);
|
||||
|
||||
actualPosition += info.size;
|
||||
}
|
||||
}
|
||||
|
||||
Close(); // We're done with this writer.
|
||||
}
|
||||
|
||||
private struct IndexInfo
|
||||
{
|
||||
public int size;
|
||||
public int typeCode;
|
||||
public uint serial;
|
||||
}
|
||||
}
|
||||
}
|
||||
/***************************************************************************
|
||||
* QueuedMemoryWriter.cs
|
||||
* -------------------
|
||||
* begin : December 16, 2010
|
||||
* copyright : (C) The RunUO Software Team
|
||||
* email : info@runuo.com
|
||||
*
|
||||
* $Id$
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
/***************************************************************************
|
||||
*
|
||||
* This program is free software; you can redistribute it and/or modify
|
||||
* it under the terms of the GNU General Public License as published by
|
||||
* the Free Software Foundation; either version 2 of the License, or
|
||||
* (at your option) any later version.
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
using System.Collections.Generic;
|
||||
using System.IO;
|
||||
|
||||
namespace Server
|
||||
{
|
||||
public sealed class QueuedMemoryWriter : BinaryFileWriter
|
||||
{
|
||||
private readonly MemoryStream _memStream;
|
||||
private readonly List<IndexInfo> _orderedIndexInfo = new List<IndexInfo>();
|
||||
|
||||
public QueuedMemoryWriter()
|
||||
: base(new MemoryStream(1024 * 1024), true) =>
|
||||
_memStream = UnderlyingStream as MemoryStream;
|
||||
|
||||
protected override int BufferSize => 512;
|
||||
|
||||
public void QueueForIndex(ISerializable serializable, int size)
|
||||
{
|
||||
IndexInfo info;
|
||||
|
||||
info.size = size;
|
||||
|
||||
info.typeCode = serializable.TypeRef; // For guilds, this will automagically be zero.
|
||||
info.serial = serializable.Serial;
|
||||
|
||||
_orderedIndexInfo.Add(info);
|
||||
}
|
||||
|
||||
public void CommitTo(SequentialFileWriterStream dataFile, SequentialFileWriterStream indexFile)
|
||||
{
|
||||
Flush();
|
||||
|
||||
var memLength = (int)_memStream.Position;
|
||||
|
||||
if (memLength > 0)
|
||||
{
|
||||
var memBuffer = _memStream.GetBuffer();
|
||||
|
||||
var actualPosition = dataFile.Position;
|
||||
|
||||
dataFile.Write(memBuffer, 0, memLength); // The buffer contains the data from many items.
|
||||
|
||||
// Console.WriteLine("Writing {0} bytes starting at {1}, with {2} things", memLength, actualPosition, _orderedIndexInfo.Count);
|
||||
|
||||
var indexBuffer = new byte[20];
|
||||
|
||||
// int indexWritten = _orderedIndexInfo.Count * indexBuffer.Length;
|
||||
// int totalWritten = memLength + indexWritten
|
||||
|
||||
for (var i = 0; i < _orderedIndexInfo.Count; i++)
|
||||
{
|
||||
var info = _orderedIndexInfo[i];
|
||||
|
||||
indexBuffer[0] = (byte)info.typeCode;
|
||||
indexBuffer[1] = (byte)(info.typeCode >> 8);
|
||||
indexBuffer[2] = (byte)(info.typeCode >> 16);
|
||||
indexBuffer[3] = (byte)(info.typeCode >> 24);
|
||||
|
||||
indexBuffer[4] = (byte)info.serial;
|
||||
indexBuffer[5] = (byte)(info.serial >> 8);
|
||||
indexBuffer[6] = (byte)(info.serial >> 16);
|
||||
indexBuffer[7] = (byte)(info.serial >> 24);
|
||||
|
||||
indexBuffer[8] = (byte)actualPosition;
|
||||
indexBuffer[9] = (byte)(actualPosition >> 8);
|
||||
indexBuffer[10] = (byte)(actualPosition >> 16);
|
||||
indexBuffer[11] = (byte)(actualPosition >> 24);
|
||||
indexBuffer[12] = (byte)(actualPosition >> 32);
|
||||
indexBuffer[13] = (byte)(actualPosition >> 40);
|
||||
indexBuffer[14] = (byte)(actualPosition >> 48);
|
||||
indexBuffer[15] = (byte)(actualPosition >> 56);
|
||||
|
||||
indexBuffer[16] = (byte)info.size;
|
||||
indexBuffer[17] = (byte)(info.size >> 8);
|
||||
indexBuffer[18] = (byte)(info.size >> 16);
|
||||
indexBuffer[19] = (byte)(info.size >> 24);
|
||||
|
||||
indexFile.Write(indexBuffer, 0, indexBuffer.Length);
|
||||
|
||||
actualPosition += info.size;
|
||||
}
|
||||
}
|
||||
|
||||
Close(); // We're done with this writer.
|
||||
}
|
||||
|
||||
private struct IndexInfo
|
||||
{
|
||||
public int size;
|
||||
public int typeCode;
|
||||
public uint serial;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,47 +1,47 @@
|
|||
/***************************************************************************
|
||||
* SaveStrategy.cs
|
||||
* -------------------
|
||||
* begin : May 1, 2002
|
||||
* copyright : (C) The RunUO Software Team
|
||||
* email : info@runuo.com
|
||||
*
|
||||
* $Id$
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
/***************************************************************************
|
||||
*
|
||||
* This program is free software; you can redistribute it and/or modify
|
||||
* it under the terms of the GNU General Public License as published by
|
||||
* the Free Software Foundation; either version 2 of the License, or
|
||||
* (at your option) any later version.
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
namespace Server
|
||||
{
|
||||
public abstract class SaveStrategy
|
||||
{
|
||||
public abstract string Name { get; }
|
||||
|
||||
public static SaveStrategy Acquire()
|
||||
{
|
||||
if (Core.MultiProcessor)
|
||||
{
|
||||
var processorCount = Core.ProcessorCount;
|
||||
|
||||
if (processorCount > 2)
|
||||
return
|
||||
new DualSaveStrategy(); // return new DynamicSaveStrategy(); (4.0 or return new ParallelSaveStrategy(processorCount); (2.0)
|
||||
|
||||
return new DualSaveStrategy();
|
||||
}
|
||||
|
||||
return new StandardSaveStrategy();
|
||||
}
|
||||
|
||||
public abstract void Save(bool permitBackgroundWrite);
|
||||
|
||||
public abstract void ProcessDecay();
|
||||
}
|
||||
}
|
||||
/***************************************************************************
|
||||
* SaveStrategy.cs
|
||||
* -------------------
|
||||
* begin : May 1, 2002
|
||||
* copyright : (C) The RunUO Software Team
|
||||
* email : info@runuo.com
|
||||
*
|
||||
* $Id$
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
/***************************************************************************
|
||||
*
|
||||
* This program is free software; you can redistribute it and/or modify
|
||||
* it under the terms of the GNU General Public License as published by
|
||||
* the Free Software Foundation; either version 2 of the License, or
|
||||
* (at your option) any later version.
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
namespace Server
|
||||
{
|
||||
public abstract class SaveStrategy
|
||||
{
|
||||
public abstract string Name { get; }
|
||||
|
||||
public static SaveStrategy Acquire()
|
||||
{
|
||||
if (Core.MultiProcessor)
|
||||
{
|
||||
var processorCount = Core.ProcessorCount;
|
||||
|
||||
if (processorCount > 2)
|
||||
return
|
||||
new DualSaveStrategy(); // return new DynamicSaveStrategy(); (4.0 or return new ParallelSaveStrategy(processorCount); (2.0)
|
||||
|
||||
return new DualSaveStrategy();
|
||||
}
|
||||
|
||||
return new StandardSaveStrategy();
|
||||
}
|
||||
|
||||
public abstract void Save(bool permitBackgroundWrite);
|
||||
|
||||
public abstract void ProcessDecay();
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,119 +1,120 @@
|
|||
/***************************************************************************
|
||||
* SequentialFileWriter.cs
|
||||
* -------------------
|
||||
* begin : May 1, 2002
|
||||
* copyright : (C) The RunUO Software Team
|
||||
* email : info@runuo.com
|
||||
*
|
||||
* $Id$
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
/***************************************************************************
|
||||
*
|
||||
* This program is free software; you can redistribute it and/or modify
|
||||
* it under the terms of the GNU General Public License as published by
|
||||
* the Free Software Foundation; either version 2 of the License, or
|
||||
* (at your option) any later version.
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
using System;
|
||||
using System.IO;
|
||||
|
||||
namespace Server
|
||||
{
|
||||
public sealed class SequentialFileWriterStream : Stream
|
||||
{
|
||||
private FileQueue fileQueue;
|
||||
private FileStream fileStream;
|
||||
|
||||
private AsyncCallback writeCallback;
|
||||
|
||||
public SequentialFileWriterStream(string path)
|
||||
{
|
||||
if (path == null) throw new ArgumentNullException(nameof(path));
|
||||
|
||||
fileStream = FileOperations.OpenSequentialStream(path, FileMode.Create, FileAccess.Write, FileShare.None);
|
||||
|
||||
fileQueue = new FileQueue(
|
||||
Math.Max(FileOperations.Concurrency, 1),
|
||||
FileCallback);
|
||||
}
|
||||
|
||||
public override long Position
|
||||
{
|
||||
get => fileQueue.Position;
|
||||
set => throw new InvalidOperationException();
|
||||
}
|
||||
|
||||
public override bool CanRead => false;
|
||||
|
||||
public override bool CanSeek => false;
|
||||
|
||||
public override bool CanWrite => true;
|
||||
|
||||
public override long Length => Position;
|
||||
|
||||
private void FileCallback(FileQueue.Chunk chunk)
|
||||
{
|
||||
if (FileOperations.AreSynchronous)
|
||||
{
|
||||
fileStream.Write(chunk.Buffer, chunk.Offset, chunk.Size);
|
||||
|
||||
chunk.Commit();
|
||||
}
|
||||
else
|
||||
{
|
||||
writeCallback ??= OnWrite;
|
||||
|
||||
fileStream.BeginWrite(chunk.Buffer, chunk.Offset, chunk.Size, writeCallback, chunk);
|
||||
}
|
||||
}
|
||||
|
||||
private void OnWrite(IAsyncResult asyncResult)
|
||||
{
|
||||
var chunk = asyncResult.AsyncState as FileQueue.Chunk;
|
||||
|
||||
fileStream.EndWrite(asyncResult);
|
||||
|
||||
chunk?.Commit();
|
||||
}
|
||||
|
||||
public override void Write(byte[] buffer, int offset, int size)
|
||||
{
|
||||
fileQueue.Enqueue(buffer, offset, size);
|
||||
}
|
||||
|
||||
public override void Flush()
|
||||
{
|
||||
fileQueue.Flush();
|
||||
fileStream.Flush();
|
||||
}
|
||||
|
||||
protected override void Dispose(bool disposing)
|
||||
{
|
||||
if (fileStream != null)
|
||||
{
|
||||
Flush();
|
||||
|
||||
fileQueue.Dispose();
|
||||
fileQueue = null;
|
||||
|
||||
fileStream.Close();
|
||||
fileStream = null;
|
||||
}
|
||||
|
||||
base.Dispose(disposing);
|
||||
}
|
||||
|
||||
public override int Read(byte[] buffer, int offset, int count) => throw new InvalidOperationException();
|
||||
|
||||
public override long Seek(long offset, SeekOrigin origin) => throw new InvalidOperationException();
|
||||
|
||||
public override void SetLength(long value)
|
||||
{
|
||||
fileStream.SetLength(value);
|
||||
}
|
||||
}
|
||||
}
|
||||
/***************************************************************************
|
||||
* SequentialFileWriter.cs
|
||||
* -------------------
|
||||
* begin : May 1, 2002
|
||||
* copyright : (C) The RunUO Software Team
|
||||
* email : info@runuo.com
|
||||
*
|
||||
* $Id$
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
/***************************************************************************
|
||||
*
|
||||
* This program is free software; you can redistribute it and/or modify
|
||||
* it under the terms of the GNU General Public License as published by
|
||||
* the Free Software Foundation; either version 2 of the License, or
|
||||
* (at your option) any later version.
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
using System;
|
||||
using System.IO;
|
||||
|
||||
namespace Server
|
||||
{
|
||||
public sealed class SequentialFileWriterStream : Stream
|
||||
{
|
||||
private FileQueue fileQueue;
|
||||
private FileStream fileStream;
|
||||
|
||||
private AsyncCallback writeCallback;
|
||||
|
||||
public SequentialFileWriterStream(string path)
|
||||
{
|
||||
if (path == null) throw new ArgumentNullException(nameof(path));
|
||||
|
||||
fileStream = FileOperations.OpenSequentialStream(path, FileMode.Create, FileAccess.Write, FileShare.None);
|
||||
|
||||
fileQueue = new FileQueue(
|
||||
Math.Max(FileOperations.Concurrency, 1),
|
||||
FileCallback
|
||||
);
|
||||
}
|
||||
|
||||
public override long Position
|
||||
{
|
||||
get => fileQueue.Position;
|
||||
set => throw new InvalidOperationException();
|
||||
}
|
||||
|
||||
public override bool CanRead => false;
|
||||
|
||||
public override bool CanSeek => false;
|
||||
|
||||
public override bool CanWrite => true;
|
||||
|
||||
public override long Length => Position;
|
||||
|
||||
private void FileCallback(FileQueue.Chunk chunk)
|
||||
{
|
||||
if (FileOperations.AreSynchronous)
|
||||
{
|
||||
fileStream.Write(chunk.Buffer, chunk.Offset, chunk.Size);
|
||||
|
||||
chunk.Commit();
|
||||
}
|
||||
else
|
||||
{
|
||||
writeCallback ??= OnWrite;
|
||||
|
||||
fileStream.BeginWrite(chunk.Buffer, chunk.Offset, chunk.Size, writeCallback, chunk);
|
||||
}
|
||||
}
|
||||
|
||||
private void OnWrite(IAsyncResult asyncResult)
|
||||
{
|
||||
var chunk = asyncResult.AsyncState as FileQueue.Chunk;
|
||||
|
||||
fileStream.EndWrite(asyncResult);
|
||||
|
||||
chunk?.Commit();
|
||||
}
|
||||
|
||||
public override void Write(byte[] buffer, int offset, int size)
|
||||
{
|
||||
fileQueue.Enqueue(buffer, offset, size);
|
||||
}
|
||||
|
||||
public override void Flush()
|
||||
{
|
||||
fileQueue.Flush();
|
||||
fileStream.Flush();
|
||||
}
|
||||
|
||||
protected override void Dispose(bool disposing)
|
||||
{
|
||||
if (fileStream != null)
|
||||
{
|
||||
Flush();
|
||||
|
||||
fileQueue.Dispose();
|
||||
fileQueue = null;
|
||||
|
||||
fileStream.Close();
|
||||
fileStream = null;
|
||||
}
|
||||
|
||||
base.Dispose(disposing);
|
||||
}
|
||||
|
||||
public override int Read(byte[] buffer, int offset, int count) => throw new InvalidOperationException();
|
||||
|
||||
public override long Seek(long offset, SeekOrigin origin) => throw new InvalidOperationException();
|
||||
|
||||
public override void SetLength(long value)
|
||||
{
|
||||
fileStream.SetLength(value);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,234 +1,250 @@
|
|||
/***************************************************************************
|
||||
* StandardSaveStrategy.cs
|
||||
* -------------------
|
||||
* begin : May 1, 2002
|
||||
* copyright : (C) The RunUO Software Team
|
||||
* email : info@runuo.com
|
||||
*
|
||||
* $Id$
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
/***************************************************************************
|
||||
*
|
||||
* This program is free software; you can redistribute it and/or modify
|
||||
* it under the terms of the GNU General Public License as published by
|
||||
* the Free Software Foundation; either version 2 of the License, or
|
||||
* (at your option) any later version.
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
using System;
|
||||
using System.Collections.Generic;
|
||||
using System.Threading.Tasks;
|
||||
using Server.Guilds;
|
||||
|
||||
namespace Server
|
||||
{
|
||||
public class StandardSaveStrategy : SaveStrategy
|
||||
{
|
||||
public enum SaveOption
|
||||
{
|
||||
Normal,
|
||||
Threaded
|
||||
}
|
||||
|
||||
// TODO: Move to configuration
|
||||
public static SaveOption SaveType => SaveOption.Normal;
|
||||
|
||||
private readonly Queue<Item> _decayQueue;
|
||||
|
||||
public StandardSaveStrategy() => _decayQueue = new Queue<Item>();
|
||||
|
||||
public override string Name => "Standard";
|
||||
|
||||
protected bool PermitBackgroundWrite { get; set; }
|
||||
|
||||
protected bool UseSequentialWriters => SaveType == SaveOption.Normal || !PermitBackgroundWrite;
|
||||
|
||||
public override void Save(bool permitBackgroundWrite)
|
||||
{
|
||||
PermitBackgroundWrite = permitBackgroundWrite;
|
||||
|
||||
#pragma warning disable CA2008 // Do not create tasks without passing a TaskScheduler *
|
||||
Task.WaitAll(Task.Factory.StartNew(SaveMobiles), Task.Factory.StartNew(SaveItems), Task.Factory.StartNew(SaveGuilds));
|
||||
#pragma warning restore CA2008 // Do not create tasks without passing a TaskScheduler *
|
||||
|
||||
if (permitBackgroundWrite && UseSequentialWriters) // If we're permitted to write in the background, but we don't anyways, then notify.
|
||||
World.NotifyDiskWriteComplete();
|
||||
}
|
||||
|
||||
protected void SaveMobiles()
|
||||
{
|
||||
var mobiles = World.Mobiles;
|
||||
|
||||
IGenericWriter idx;
|
||||
IGenericWriter tdb;
|
||||
IGenericWriter bin;
|
||||
|
||||
if (UseSequentialWriters)
|
||||
{
|
||||
idx = new BinaryFileWriter(World.MobileIndexPath, false);
|
||||
tdb = new BinaryFileWriter(World.MobileTypesPath, false);
|
||||
bin = new BinaryFileWriter(World.MobileDataPath, true);
|
||||
}
|
||||
else
|
||||
{
|
||||
idx = new AsyncWriter(World.MobileIndexPath, false);
|
||||
tdb = new AsyncWriter(World.MobileTypesPath, false);
|
||||
bin = new AsyncWriter(World.MobileDataPath, true);
|
||||
}
|
||||
#pragma warning disable CA2008 // Do not create tasks without passing a TaskScheduler *
|
||||
Task.Factory.StartNew(() =>
|
||||
{
|
||||
tdb.Write(World.m_MobileTypes.Count);
|
||||
|
||||
for (var i = 0; i < World.m_MobileTypes.Count; ++i)
|
||||
tdb.Write(World.m_MobileTypes[i].FullName);
|
||||
|
||||
tdb.Close();
|
||||
});
|
||||
#pragma warning restore CA2008 // Do not create tasks without passing a TaskScheduler *
|
||||
|
||||
Parallel.ForEach(mobiles.Values, mobile => mobile.Serialize());
|
||||
#pragma warning disable CA2008 // Do not create tasks without passing a TaskScheduler *
|
||||
Task.Factory.StartNew(() =>
|
||||
{
|
||||
idx.Write(mobiles.Count);
|
||||
foreach (var m in mobiles.Values)
|
||||
{
|
||||
var start = bin.Position;
|
||||
|
||||
idx.Write(m.TypeRef);
|
||||
idx.Write(m.Serial);
|
||||
idx.Write(start);
|
||||
idx.Write((int)m.SaveBuffer.Position);
|
||||
|
||||
m.SaveBuffer.WriteTo(bin);
|
||||
m.FreeCache();
|
||||
}
|
||||
|
||||
idx.Close();
|
||||
bin.Close();
|
||||
});
|
||||
#pragma warning restore CA2008 // Do not create tasks without passing a TaskScheduler *
|
||||
}
|
||||
|
||||
protected void SaveItems()
|
||||
{
|
||||
var items = World.Items;
|
||||
|
||||
IGenericWriter idx;
|
||||
IGenericWriter tdb;
|
||||
IGenericWriter bin;
|
||||
|
||||
if (UseSequentialWriters)
|
||||
{
|
||||
idx = new BinaryFileWriter(World.ItemIndexPath, false);
|
||||
tdb = new BinaryFileWriter(World.ItemTypesPath, false);
|
||||
bin = new BinaryFileWriter(World.ItemDataPath, true);
|
||||
}
|
||||
else
|
||||
{
|
||||
idx = new AsyncWriter(World.ItemIndexPath, false);
|
||||
tdb = new AsyncWriter(World.ItemTypesPath, false);
|
||||
bin = new AsyncWriter(World.ItemDataPath, true);
|
||||
}
|
||||
|
||||
#pragma warning disable CA2008 // Do not create tasks without passing a TaskScheduler *
|
||||
Task.Factory.StartNew(() =>
|
||||
{
|
||||
tdb.Write(World.m_ItemTypes.Count);
|
||||
|
||||
for (var i = 0; i < World.m_ItemTypes.Count; ++i)
|
||||
tdb.Write(World.m_ItemTypes[i].FullName);
|
||||
|
||||
tdb.Close();
|
||||
});
|
||||
#pragma warning restore CA2008 // Do not create tasks without passing a TaskScheduler *
|
||||
Parallel.ForEach(items.Values, item => item.Serialize());
|
||||
|
||||
idx.Write(items.Count);
|
||||
|
||||
var n = DateTime.UtcNow;
|
||||
|
||||
#pragma warning disable CA2008 // Do not create tasks without passing a TaskScheduler *
|
||||
Task.Factory.StartNew(() =>
|
||||
{
|
||||
foreach (var item in items.Values)
|
||||
{
|
||||
if (item.Decays && item.Parent == null && item.Map != Map.Internal && item.LastMoved + item.DecayTime <= n)
|
||||
{
|
||||
Console.WriteLine($"Decay Item {item.Name ?? item.DefaultName} ({item.GetType().FullName})");
|
||||
_decayQueue.Enqueue(item);
|
||||
}
|
||||
|
||||
var start = bin.Position;
|
||||
|
||||
idx.Write(item.TypeRef);
|
||||
idx.Write(item.Serial);
|
||||
idx.Write(start);
|
||||
idx.Write((int)item.SaveBuffer.Position);
|
||||
|
||||
item.SaveBuffer.WriteTo(bin);
|
||||
item.FreeCache();
|
||||
}
|
||||
|
||||
idx.Close();
|
||||
bin.Close();
|
||||
});
|
||||
#pragma warning restore CA2008 // Do not create tasks without passing a TaskScheduler *
|
||||
}
|
||||
|
||||
protected void SaveGuilds()
|
||||
{
|
||||
IGenericWriter idx;
|
||||
IGenericWriter bin;
|
||||
|
||||
if (UseSequentialWriters)
|
||||
{
|
||||
idx = new BinaryFileWriter(World.GuildIndexPath, false);
|
||||
bin = new BinaryFileWriter(World.GuildDataPath, true);
|
||||
}
|
||||
else
|
||||
{
|
||||
idx = new AsyncWriter(World.GuildIndexPath, false);
|
||||
bin = new AsyncWriter(World.GuildDataPath, true);
|
||||
}
|
||||
|
||||
Parallel.ForEach(BaseGuild.List.Values, guild => guild.Serialize());
|
||||
|
||||
idx.Write(BaseGuild.List.Count);
|
||||
#pragma warning disable CA2008 // Do not create tasks without passing a TaskScheduler *
|
||||
Task.Factory.StartNew(() =>
|
||||
{
|
||||
foreach (var guild in BaseGuild.List.Values)
|
||||
{
|
||||
var start = bin.Position;
|
||||
|
||||
idx.Write(0); // guilds have no typeid
|
||||
idx.Write(guild.Serial);
|
||||
idx.Write(start);
|
||||
idx.Write((int)guild.SaveBuffer.Position);
|
||||
|
||||
guild.SaveBuffer.WriteTo(bin);
|
||||
}
|
||||
|
||||
idx.Close();
|
||||
bin.Close();
|
||||
});
|
||||
#pragma warning restore CA2008 // Do not create tasks without passing a TaskScheduler *
|
||||
}
|
||||
|
||||
public override void ProcessDecay()
|
||||
{
|
||||
while (_decayQueue.Count > 0)
|
||||
{
|
||||
var item = _decayQueue.Dequeue();
|
||||
|
||||
if (item.OnDecay())
|
||||
item.Delete();
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
/***************************************************************************
|
||||
* StandardSaveStrategy.cs
|
||||
* -------------------
|
||||
* begin : May 1, 2002
|
||||
* copyright : (C) The RunUO Software Team
|
||||
* email : info@runuo.com
|
||||
*
|
||||
* $Id$
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
/***************************************************************************
|
||||
*
|
||||
* This program is free software; you can redistribute it and/or modify
|
||||
* it under the terms of the GNU General Public License as published by
|
||||
* the Free Software Foundation; either version 2 of the License, or
|
||||
* (at your option) any later version.
|
||||
*
|
||||
***************************************************************************/
|
||||
|
||||
using System;
|
||||
using System.Collections.Generic;
|
||||
using System.Threading.Tasks;
|
||||
using Server.Guilds;
|
||||
|
||||
namespace Server
|
||||
{
|
||||
public class StandardSaveStrategy : SaveStrategy
|
||||
{
|
||||
public enum SaveOption
|
||||
{
|
||||
Normal,
|
||||
Threaded
|
||||
}
|
||||
|
||||
private readonly Queue<Item> _decayQueue;
|
||||
|
||||
public StandardSaveStrategy() => _decayQueue = new Queue<Item>();
|
||||
|
||||
// TODO: Move to configuration
|
||||
public static SaveOption SaveType => SaveOption.Normal;
|
||||
|
||||
public override string Name => "Standard";
|
||||
|
||||
protected bool PermitBackgroundWrite { get; set; }
|
||||
|
||||
protected bool UseSequentialWriters => SaveType == SaveOption.Normal || !PermitBackgroundWrite;
|
||||
|
||||
public override void Save(bool permitBackgroundWrite)
|
||||
{
|
||||
PermitBackgroundWrite = permitBackgroundWrite;
|
||||
|
||||
#pragma warning disable CA2008 // Do not create tasks without passing a TaskScheduler *
|
||||
Task.WaitAll(
|
||||
Task.Factory.StartNew(SaveMobiles),
|
||||
Task.Factory.StartNew(SaveItems),
|
||||
Task.Factory.StartNew(SaveGuilds)
|
||||
);
|
||||
#pragma warning restore CA2008 // Do not create tasks without passing a TaskScheduler *
|
||||
|
||||
if (permitBackgroundWrite && UseSequentialWriters
|
||||
) // If we're permitted to write in the background, but we don't anyways, then notify.
|
||||
World.NotifyDiskWriteComplete();
|
||||
}
|
||||
|
||||
protected void SaveMobiles()
|
||||
{
|
||||
var mobiles = World.Mobiles;
|
||||
|
||||
IGenericWriter idx;
|
||||
IGenericWriter tdb;
|
||||
IGenericWriter bin;
|
||||
|
||||
if (UseSequentialWriters)
|
||||
{
|
||||
idx = new BinaryFileWriter(World.MobileIndexPath, false);
|
||||
tdb = new BinaryFileWriter(World.MobileTypesPath, false);
|
||||
bin = new BinaryFileWriter(World.MobileDataPath, true);
|
||||
}
|
||||
else
|
||||
{
|
||||
idx = new AsyncWriter(World.MobileIndexPath, false);
|
||||
tdb = new AsyncWriter(World.MobileTypesPath, false);
|
||||
bin = new AsyncWriter(World.MobileDataPath, true);
|
||||
}
|
||||
#pragma warning disable CA2008 // Do not create tasks without passing a TaskScheduler *
|
||||
Task.Factory.StartNew(
|
||||
() =>
|
||||
{
|
||||
tdb.Write(World.m_MobileTypes.Count);
|
||||
|
||||
for (var i = 0; i < World.m_MobileTypes.Count; ++i)
|
||||
tdb.Write(World.m_MobileTypes[i].FullName);
|
||||
|
||||
tdb.Close();
|
||||
}
|
||||
);
|
||||
#pragma warning restore CA2008 // Do not create tasks without passing a TaskScheduler *
|
||||
|
||||
Parallel.ForEach(mobiles.Values, mobile => mobile.Serialize());
|
||||
#pragma warning disable CA2008 // Do not create tasks without passing a TaskScheduler *
|
||||
Task.Factory.StartNew(
|
||||
() =>
|
||||
{
|
||||
idx.Write(mobiles.Count);
|
||||
foreach (var m in mobiles.Values)
|
||||
{
|
||||
var start = bin.Position;
|
||||
|
||||
idx.Write(m.TypeRef);
|
||||
idx.Write(m.Serial);
|
||||
idx.Write(start);
|
||||
idx.Write((int)m.SaveBuffer.Position);
|
||||
|
||||
m.SaveBuffer.WriteTo(bin);
|
||||
m.FreeCache();
|
||||
}
|
||||
|
||||
idx.Close();
|
||||
bin.Close();
|
||||
}
|
||||
);
|
||||
#pragma warning restore CA2008 // Do not create tasks without passing a TaskScheduler *
|
||||
}
|
||||
|
||||
protected void SaveItems()
|
||||
{
|
||||
var items = World.Items;
|
||||
|
||||
IGenericWriter idx;
|
||||
IGenericWriter tdb;
|
||||
IGenericWriter bin;
|
||||
|
||||
if (UseSequentialWriters)
|
||||
{
|
||||
idx = new BinaryFileWriter(World.ItemIndexPath, false);
|
||||
tdb = new BinaryFileWriter(World.ItemTypesPath, false);
|
||||
bin = new BinaryFileWriter(World.ItemDataPath, true);
|
||||
}
|
||||
else
|
||||
{
|
||||
idx = new AsyncWriter(World.ItemIndexPath, false);
|
||||
tdb = new AsyncWriter(World.ItemTypesPath, false);
|
||||
bin = new AsyncWriter(World.ItemDataPath, true);
|
||||
}
|
||||
|
||||
#pragma warning disable CA2008 // Do not create tasks without passing a TaskScheduler *
|
||||
Task.Factory.StartNew(
|
||||
() =>
|
||||
{
|
||||
tdb.Write(World.m_ItemTypes.Count);
|
||||
|
||||
for (var i = 0; i < World.m_ItemTypes.Count; ++i)
|
||||
tdb.Write(World.m_ItemTypes[i].FullName);
|
||||
|
||||
tdb.Close();
|
||||
}
|
||||
);
|
||||
#pragma warning restore CA2008 // Do not create tasks without passing a TaskScheduler *
|
||||
Parallel.ForEach(items.Values, item => item.Serialize());
|
||||
|
||||
idx.Write(items.Count);
|
||||
|
||||
var n = DateTime.UtcNow;
|
||||
|
||||
#pragma warning disable CA2008 // Do not create tasks without passing a TaskScheduler *
|
||||
Task.Factory.StartNew(
|
||||
() =>
|
||||
{
|
||||
foreach (var item in items.Values)
|
||||
{
|
||||
if (item.Decays && item.Parent == null && item.Map != Map.Internal &&
|
||||
item.LastMoved + item.DecayTime <= n)
|
||||
{
|
||||
Console.WriteLine($"Decay Item {item.Name ?? item.DefaultName} ({item.GetType().FullName})");
|
||||
_decayQueue.Enqueue(item);
|
||||
}
|
||||
|
||||
var start = bin.Position;
|
||||
|
||||
idx.Write(item.TypeRef);
|
||||
idx.Write(item.Serial);
|
||||
idx.Write(start);
|
||||
idx.Write((int)item.SaveBuffer.Position);
|
||||
|
||||
item.SaveBuffer.WriteTo(bin);
|
||||
item.FreeCache();
|
||||
}
|
||||
|
||||
idx.Close();
|
||||
bin.Close();
|
||||
}
|
||||
);
|
||||
#pragma warning restore CA2008 // Do not create tasks without passing a TaskScheduler *
|
||||
}
|
||||
|
||||
protected void SaveGuilds()
|
||||
{
|
||||
IGenericWriter idx;
|
||||
IGenericWriter bin;
|
||||
|
||||
if (UseSequentialWriters)
|
||||
{
|
||||
idx = new BinaryFileWriter(World.GuildIndexPath, false);
|
||||
bin = new BinaryFileWriter(World.GuildDataPath, true);
|
||||
}
|
||||
else
|
||||
{
|
||||
idx = new AsyncWriter(World.GuildIndexPath, false);
|
||||
bin = new AsyncWriter(World.GuildDataPath, true);
|
||||
}
|
||||
|
||||
Parallel.ForEach(BaseGuild.List.Values, guild => guild.Serialize());
|
||||
|
||||
idx.Write(BaseGuild.List.Count);
|
||||
#pragma warning disable CA2008 // Do not create tasks without passing a TaskScheduler *
|
||||
Task.Factory.StartNew(
|
||||
() =>
|
||||
{
|
||||
foreach (var guild in BaseGuild.List.Values)
|
||||
{
|
||||
var start = bin.Position;
|
||||
|
||||
idx.Write(0); // guilds have no typeid
|
||||
idx.Write(guild.Serial);
|
||||
idx.Write(start);
|
||||
idx.Write((int)guild.SaveBuffer.Position);
|
||||
|
||||
guild.SaveBuffer.WriteTo(bin);
|
||||
}
|
||||
|
||||
idx.Close();
|
||||
bin.Close();
|
||||
}
|
||||
);
|
||||
#pragma warning restore CA2008 // Do not create tasks without passing a TaskScheduler *
|
||||
}
|
||||
|
||||
public override void ProcessDecay()
|
||||
{
|
||||
while (_decayQueue.Count > 0)
|
||||
{
|
||||
var item = _decayQueue.Dequeue();
|
||||
|
||||
if (item.OnDecay())
|
||||
item.Delete();
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue