Compare commits
140
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
a41272f160 | ||
|
|
8e018cbae9 | ||
|
|
3c447eb4c5 | ||
|
|
d53b760731 | ||
|
|
fcadceb009 | ||
|
|
b3013a9b65 | ||
|
|
eac06e8b3a | ||
|
|
f7fa35c7ba | ||
|
|
a7526aaed1 | ||
|
|
e7d9e833f1 | ||
|
|
13dd0a649c | ||
|
|
bf5bd8d726 | ||
|
|
8dcf9ff15e | ||
|
|
489209d549 | ||
|
|
017cee43c0 | ||
|
|
8fbef6ff51 | ||
|
|
38ee2e4eb5 | ||
|
|
e4f7249f1e | ||
|
|
e3ca516b97 | ||
|
|
88daab379f | ||
|
|
00ed3caafe | ||
|
|
daad6468c6 | ||
|
|
3d347a936d | ||
|
|
87bda475df | ||
|
|
44dfa39737 | ||
|
|
d38f6da26f | ||
|
|
5667cfc054 | ||
|
|
ecd0e70261 | ||
|
|
5bbb95f1ff | ||
|
|
5aa3edbda2 | ||
|
|
9bdebb963b | ||
|
|
c690598868 | ||
|
|
ba59ac91d7 | ||
|
|
5170699ae4 | ||
|
|
09d2f418eb | ||
|
|
cc9d04acd7 | ||
|
|
b5c8c55c72 | ||
|
|
38b987144e | ||
|
|
8136bf1c92 | ||
|
|
5ddec99114 | ||
|
|
58cebd43ce | ||
|
|
11a015e0b6 | ||
|
|
745b7e43b4 | ||
|
|
a72866c725 | ||
|
|
2ec7f26387 | ||
|
|
72865e7f96 | ||
|
|
6a0d2830d6 | ||
|
|
b3710f26ae | ||
|
|
f6cd9db408 | ||
|
|
7a3a2b558b | ||
|
|
53aad5cb37 | ||
|
|
fd65e1f022 | ||
|
|
d16b8cb011 | ||
|
|
37611bdc66 | ||
|
|
7352cdf3fe | ||
|
|
8e084d6ca0 | ||
|
|
017c71a6e8 | ||
|
|
488fbb383b | ||
|
|
cfd635146b | ||
|
|
7d737a5534 | ||
|
|
0d21246b06 | ||
|
|
ec99f5f8aa | ||
|
|
9d9f65c5a3 | ||
|
|
87271f4f37 | ||
|
|
9bebc23f6a | ||
|
|
682c267b10 | ||
|
|
4de35b579f | ||
|
|
fb4d520ba4 | ||
|
|
1f72d4f37d | ||
|
|
ecada92dc5 | ||
|
|
410d62849a | ||
|
|
927ccaa119 | ||
|
|
d58d65b751 | ||
|
|
802b18e990 | ||
|
|
8064051f2f | ||
|
|
485b2387e4 | ||
|
|
ae4d782566 | ||
|
|
8f16699bab | ||
|
|
5ba919a638 | ||
|
|
af7aeb2bba | ||
|
|
1f11c5d6b5 | ||
|
|
8694bfe9be | ||
|
|
637b5a4736 | ||
|
|
e0cc46260b | ||
|
|
c1ad97d26f | ||
|
|
1e8f8106b0 | ||
|
|
926e323569 | ||
|
|
8d32758918 | ||
|
|
7eed6160eb | ||
|
|
c57dc4daa1 | ||
|
|
4c75cebcd6 | ||
|
|
a820788c7d | ||
|
|
52a02abd3f | ||
|
|
a92455b2a5 | ||
|
|
933d5e7d4d | ||
|
|
ba43fd90c6 | ||
|
|
01ee514c73 | ||
|
|
1eb30329c6 | ||
|
|
0ef55abdf4 | ||
|
|
8341807d92 | ||
|
|
c2712daccb | ||
|
|
8033da1176 | ||
|
|
5afab577a7 | ||
|
|
4be9b9df9a | ||
|
|
a905f0ce53 | ||
|
|
0c4d3be912 | ||
|
|
2b7ba61543 | ||
|
|
f9c7e18985 | ||
|
|
32e7028029 | ||
|
|
cc3eddf02d | ||
|
|
d73ad6db15 | ||
|
|
f7c45d17d7 | ||
|
|
cb4cdfe69a | ||
|
|
bed57dd35b | ||
|
|
c38a2242ba | ||
|
|
d3488dc907 | ||
|
|
1fa7787d3b | ||
|
|
1c856f7615 | ||
|
|
cc2513bd2f | ||
|
|
ac7d323201 | ||
|
|
a846d51c0c | ||
|
|
f5da80dc9c | ||
|
|
f94a67adb4 | ||
|
|
f1d453251c | ||
|
|
12dc7efd5b | ||
|
|
117a30bb82 | ||
|
|
ccc6c815e4 | ||
|
|
ba9e4b098f | ||
|
|
69aa3a998f | ||
|
|
00fd2cebf9 | ||
|
|
5143361dcd | ||
|
|
b1818400ca | ||
|
|
a236544ee9 | ||
|
|
fa1b560a91 | ||
|
|
d6f7e225be | ||
|
|
22cf82b99b | ||
|
|
362040bd3c | ||
|
|
40393c3a5b | ||
|
|
7f972bac85 | ||
|
|
d532d9505a |
@@ -0,0 +1,30 @@
|
||||
**/.classpath
|
||||
**/.dockerignore
|
||||
**/.env
|
||||
**/.git
|
||||
**/.gitignore
|
||||
**/.project
|
||||
**/.settings
|
||||
**/.toolstarget
|
||||
**/.vs
|
||||
**/.vscode
|
||||
**/*.*proj.user
|
||||
**/*.dbmdl
|
||||
**/*.jfm
|
||||
**/azds.yaml
|
||||
**/bin
|
||||
**/charts
|
||||
**/docker-compose*
|
||||
**/Dockerfile*
|
||||
**/node_modules
|
||||
**/npm-debug.log
|
||||
**/obj
|
||||
**/secrets.dev.yaml
|
||||
**/values.dev.yaml
|
||||
LICENSE
|
||||
README.md
|
||||
!**/.gitignore
|
||||
!.git/HEAD
|
||||
!.git/config
|
||||
!.git/packed-refs
|
||||
!.git/refs/heads/**
|
||||
@@ -0,0 +1,26 @@
|
||||
name: Docker - AutoClient
|
||||
|
||||
on:
|
||||
push:
|
||||
branches:
|
||||
- master
|
||||
tags:
|
||||
- 'v*.*.*'
|
||||
paths:
|
||||
- 'Tools/AutoClient/**'
|
||||
- '!Tools/AutoClient/docker/docker-compose.yaml'
|
||||
- 'Framework/**'
|
||||
- 'ProjectPlugins/**'
|
||||
- .github/workflows/docker-autoclient.yml
|
||||
- .github/workflows/docker-reusable.yml
|
||||
workflow_dispatch:
|
||||
|
||||
jobs:
|
||||
build-and-push:
|
||||
name: Build and Push
|
||||
uses: ./.github/workflows/docker-reusable.yml
|
||||
with:
|
||||
docker_file: Tools/AutoClient/docker/Dockerfile
|
||||
docker_repo: codexstorage/codex-autoclient
|
||||
secrets: inherit
|
||||
|
||||
@@ -0,0 +1,27 @@
|
||||
name: Docker - MarketInsights API
|
||||
|
||||
|
||||
on:
|
||||
push:
|
||||
branches:
|
||||
- master
|
||||
tags:
|
||||
- 'v*.*.*'
|
||||
paths:
|
||||
- 'Tools/MarketInsights/**'
|
||||
- 'Framework/**'
|
||||
- 'ProjectPlugins/**'
|
||||
- .github/workflows/docker-marketinsights.yml
|
||||
- .github/workflows/docker-reusable.yml
|
||||
workflow_dispatch:
|
||||
|
||||
|
||||
jobs:
|
||||
build-and-push:
|
||||
name: Build and Push
|
||||
uses: ./.github/workflows/docker-reusable.yml
|
||||
with:
|
||||
docker_file: Tools/MarketInsights/Dockerfile
|
||||
docker_repo: codexstorage/codex-marketinsights
|
||||
secrets: inherit
|
||||
|
||||
+2
-1
@@ -1,4 +1,5 @@
|
||||
.vs
|
||||
obj
|
||||
bin
|
||||
.vscode
|
||||
.vscode
|
||||
Tools/AutoClient/datapath
|
||||
|
||||
@@ -30,11 +30,7 @@ namespace Core
|
||||
public IDownloadedLog DownloadLog(RunningContainer container, int? tailLines = null)
|
||||
{
|
||||
var workflow = entryPoint.Tools.CreateWorkflow();
|
||||
var msg = $"Downloading container log for '{container.Name}'";
|
||||
entryPoint.Tools.GetLog().Log(msg);
|
||||
var logHandler = new WriteToFileLogHandler(entryPoint.Tools.GetLog(), msg);
|
||||
workflow.DownloadContainerLog(container, logHandler, tailLines);
|
||||
return new DownloadedLog(logHandler);
|
||||
return workflow.DownloadContainerLog(container, tailLines);
|
||||
}
|
||||
|
||||
public string ExecuteContainerCommand(IHasContainer containerSource, string command, params string[] args)
|
||||
|
||||
+17
-4
@@ -13,18 +13,21 @@ namespace Core
|
||||
|
||||
internal class Http : IHttp
|
||||
{
|
||||
private static readonly object httpLock = new object();
|
||||
private static object lockLock = new object();
|
||||
private static readonly Dictionary<string, object> httpLocks = new Dictionary<string, object>();
|
||||
private readonly ILog log;
|
||||
private readonly ITimeSet timeSet;
|
||||
private readonly Action<HttpClient> onClientCreated;
|
||||
private readonly string id;
|
||||
|
||||
internal Http(ILog log, ITimeSet timeSet)
|
||||
: this(log, timeSet, DoNothing)
|
||||
internal Http(string id, ILog log, ITimeSet timeSet)
|
||||
: this(id, log, timeSet, DoNothing)
|
||||
{
|
||||
}
|
||||
|
||||
internal Http(ILog log, ITimeSet timeSet, Action<HttpClient> onClientCreated)
|
||||
internal Http(string id, ILog log, ITimeSet timeSet, Action<HttpClient> onClientCreated)
|
||||
{
|
||||
this.id = id;
|
||||
this.log = log;
|
||||
this.timeSet = timeSet;
|
||||
this.onClientCreated = onClientCreated;
|
||||
@@ -63,12 +66,22 @@ namespace Core
|
||||
|
||||
private T LockRetry<T>(Func<T> operation, Retry retry)
|
||||
{
|
||||
var httpLock = GetLock();
|
||||
lock (httpLock)
|
||||
{
|
||||
return retry.Run(operation);
|
||||
}
|
||||
}
|
||||
|
||||
private object GetLock()
|
||||
{
|
||||
lock (lockLock) // I had to.
|
||||
{
|
||||
if (!httpLocks.ContainsKey(id)) httpLocks.Add(id, new object());
|
||||
return httpLocks[id];
|
||||
}
|
||||
}
|
||||
|
||||
private HttpClient GetClient()
|
||||
{
|
||||
var client = new HttpClient();
|
||||
|
||||
@@ -27,9 +27,9 @@ namespace Core
|
||||
|
||||
public interface IHttpFactoryTool
|
||||
{
|
||||
IHttp CreateHttp(Action<HttpClient> onClientCreated);
|
||||
IHttp CreateHttp(Action<HttpClient> onClientCreated, ITimeSet timeSet);
|
||||
IHttp CreateHttp();
|
||||
IHttp CreateHttp(string id, Action<HttpClient> onClientCreated);
|
||||
IHttp CreateHttp(string id, Action<HttpClient> onClientCreated, ITimeSet timeSet);
|
||||
IHttp CreateHttp(string id);
|
||||
}
|
||||
|
||||
public interface IFileTool
|
||||
@@ -58,19 +58,19 @@ namespace Core
|
||||
log.Prefix = prefix;
|
||||
}
|
||||
|
||||
public IHttp CreateHttp(Action<HttpClient> onClientCreated)
|
||||
public IHttp CreateHttp(string id, Action<HttpClient> onClientCreated)
|
||||
{
|
||||
return CreateHttp(onClientCreated, TimeSet);
|
||||
return CreateHttp(id, onClientCreated, TimeSet);
|
||||
}
|
||||
|
||||
public IHttp CreateHttp(Action<HttpClient> onClientCreated, ITimeSet ts)
|
||||
public IHttp CreateHttp(string id, Action<HttpClient> onClientCreated, ITimeSet ts)
|
||||
{
|
||||
return new Http(log, ts, onClientCreated);
|
||||
return new Http(id, log, ts, onClientCreated);
|
||||
}
|
||||
|
||||
public IHttp CreateHttp()
|
||||
public IHttp CreateHttp(string id)
|
||||
{
|
||||
return new Http(log, TimeSet);
|
||||
return new Http(id, log, TimeSet);
|
||||
}
|
||||
|
||||
public IStartupWorkflow CreateWorkflow(string? namespaceOverride = null)
|
||||
|
||||
@@ -34,12 +34,12 @@
|
||||
{
|
||||
public TimeSpan HttpCallTimeout()
|
||||
{
|
||||
return TimeSpan.FromMinutes(3);
|
||||
return TimeSpan.FromMinutes(2);
|
||||
}
|
||||
|
||||
public TimeSpan HttpRetryTimeout()
|
||||
{
|
||||
return TimeSpan.FromMinutes(10);
|
||||
return TimeSpan.FromMinutes(5);
|
||||
}
|
||||
|
||||
public TimeSpan HttpCallRetryDelay()
|
||||
|
||||
@@ -13,9 +13,9 @@ namespace DiscordRewards
|
||||
public enum CheckType
|
||||
{
|
||||
Uninitialized,
|
||||
FilledSlot,
|
||||
FinishedSlot,
|
||||
PostedContract,
|
||||
StartedContract,
|
||||
HostFilledSlot,
|
||||
HostFinishedSlot,
|
||||
ClientPostedContract,
|
||||
ClientStartedContract,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -3,8 +3,12 @@
|
||||
public class GiveRewardsCommand
|
||||
{
|
||||
public RewardUsersCommand[] Rewards { get; set; } = Array.Empty<RewardUsersCommand>();
|
||||
public MarketAverage[] Averages { get; set; } = Array.Empty<MarketAverage>();
|
||||
public string[] EventsOverview { get; set; } = Array.Empty<string>();
|
||||
|
||||
public bool HasAny()
|
||||
{
|
||||
return Rewards.Any() || EventsOverview.Any();
|
||||
}
|
||||
}
|
||||
|
||||
public class RewardUsersCommand
|
||||
@@ -12,15 +16,4 @@
|
||||
public ulong RewardId { get; set; }
|
||||
public string[] UserAddresses { get; set; } = Array.Empty<string>();
|
||||
}
|
||||
|
||||
public class MarketAverage
|
||||
{
|
||||
public int NumberOfFinished { get; set; }
|
||||
public int TimeRangeSeconds { get; set; }
|
||||
public float Price { get; set; }
|
||||
public float Size { get; set; }
|
||||
public float Duration { get; set; }
|
||||
public float Collateral { get; set; }
|
||||
public float ProofProbability { get; set; }
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,53 +1,53 @@
|
||||
using Utils;
|
||||
|
||||
namespace DiscordRewards
|
||||
namespace DiscordRewards
|
||||
{
|
||||
public class RewardRepo
|
||||
{
|
||||
private static string Tag => RewardConfig.UsernameTag;
|
||||
|
||||
public RewardConfig[] Rewards { get; } = new RewardConfig[]
|
||||
{
|
||||
// Filled any slot
|
||||
new RewardConfig(1187039439558541498, $"{Tag} successfully filled their first slot!", new CheckConfig
|
||||
{
|
||||
Type = CheckType.FilledSlot
|
||||
}),
|
||||
public RewardConfig[] Rewards { get; } = new RewardConfig[0];
|
||||
|
||||
// Finished any slot
|
||||
new RewardConfig(1202286165630390339, $"{Tag} successfully finished their first slot!", new CheckConfig
|
||||
{
|
||||
Type = CheckType.FinishedSlot
|
||||
}),
|
||||
// Example configuration, from test server:
|
||||
//{
|
||||
// // Filled any slot
|
||||
// new RewardConfig(1187039439558541498, $"{Tag} successfully filled their first slot!", new CheckConfig
|
||||
// {
|
||||
// Type = CheckType.HostFilledSlot
|
||||
// }),
|
||||
|
||||
// Finished a sizable slot
|
||||
new RewardConfig(1202286218738405418, $"{Tag} finished their first 1GB-24h slot! (10mb/5mins for test)", new CheckConfig
|
||||
{
|
||||
Type = CheckType.FinishedSlot,
|
||||
MinSlotSize = 10.MB(),
|
||||
MinDuration = TimeSpan.FromMinutes(5.0),
|
||||
}),
|
||||
// // Finished any slot
|
||||
// new RewardConfig(1202286165630390339, $"{Tag} successfully finished their first slot!", new CheckConfig
|
||||
// {
|
||||
// Type = CheckType.HostFinishedSlot
|
||||
// }),
|
||||
|
||||
// Posted any contract
|
||||
new RewardConfig(1202286258370383913, $"{Tag} posted their first contract!", new CheckConfig
|
||||
{
|
||||
Type = CheckType.PostedContract
|
||||
}),
|
||||
// // Finished a sizable slot
|
||||
// new RewardConfig(1202286218738405418, $"{Tag} finished their first 1GB-24h slot! (10mb/5mins for test)", new CheckConfig
|
||||
// {
|
||||
// Type = CheckType.HostFinishedSlot,
|
||||
// MinSlotSize = 10.MB(),
|
||||
// MinDuration = TimeSpan.FromMinutes(5.0),
|
||||
// }),
|
||||
|
||||
// Started any contract
|
||||
new RewardConfig(1202286330873126992, $"A contract created by {Tag} reached Started state for the first time!", new CheckConfig
|
||||
{
|
||||
Type = CheckType.StartedContract
|
||||
}),
|
||||
// // Posted any contract
|
||||
// new RewardConfig(1202286258370383913, $"{Tag} posted their first contract!", new CheckConfig
|
||||
// {
|
||||
// Type = CheckType.ClientPostedContract
|
||||
// }),
|
||||
|
||||
// Started a sizable contract
|
||||
new RewardConfig(1202286381670608909, $"A large contract created by {Tag} reached Started state for the first time! (10mb/5mins for test)", new CheckConfig
|
||||
{
|
||||
Type = CheckType.StartedContract,
|
||||
MinNumberOfHosts = 4,
|
||||
MinSlotSize = 10.MB(),
|
||||
MinDuration = TimeSpan.FromMinutes(5.0),
|
||||
})
|
||||
};
|
||||
// // Started any contract
|
||||
// new RewardConfig(1202286330873126992, $"A contract created by {Tag} reached Started state for the first time!", new CheckConfig
|
||||
// {
|
||||
// Type = CheckType.ClientStartedContract
|
||||
// }),
|
||||
|
||||
// // Started a sizable contract
|
||||
// new RewardConfig(1202286381670608909, $"A large contract created by {Tag} reached Started state for the first time! (10mb/5mins for test)", new CheckConfig
|
||||
// {
|
||||
// Type = CheckType.ClientStartedContract,
|
||||
// MinNumberOfHosts = 4,
|
||||
// MinSlotSize = 10.MB(),
|
||||
// MinDuration = TimeSpan.FromMinutes(5.0),
|
||||
// })
|
||||
//};
|
||||
}
|
||||
}
|
||||
|
||||
@@ -7,27 +7,35 @@ namespace FileUtils
|
||||
{
|
||||
TrackedFile CreateEmptyFile(string label = "");
|
||||
TrackedFile GenerateFile(ByteSize size, string label = "");
|
||||
TrackedFile GenerateFile(Action<IGenerateOption> options, string label = "");
|
||||
void DeleteAllFiles();
|
||||
void ScopedFiles(Action action);
|
||||
T ScopedFiles<T>(Func<T> action);
|
||||
}
|
||||
|
||||
public interface IGenerateOption
|
||||
{
|
||||
IGenerateOption Random(ByteSize size);
|
||||
IGenerateOption StringRepeat(string str, ByteSize size);
|
||||
IGenerateOption StringRepeat(string str, int times);
|
||||
IGenerateOption ByteRepeat(byte[] bytes, ByteSize size);
|
||||
IGenerateOption ByteRepeat(byte[] bytes, int times);
|
||||
}
|
||||
|
||||
public class FileManager : IFileManager
|
||||
{
|
||||
public const int ChunkSize = 1024 * 1024 * 100;
|
||||
private static NumberSource folderNumberSource = new NumberSource(0);
|
||||
private readonly Random random = new Random();
|
||||
private static readonly NumberSource folderNumberSource = new NumberSource(0);
|
||||
private readonly ILog log;
|
||||
private readonly string rootFolder;
|
||||
private readonly string folder;
|
||||
private readonly List<List<TrackedFile>> fileSetStack = new List<List<TrackedFile>>();
|
||||
|
||||
public const int ChunkSize = 1024 * 1024 * 100;
|
||||
|
||||
public FileManager(ILog log, string rootFolder)
|
||||
{
|
||||
folder = Path.Combine(rootFolder, folderNumberSource.GetNextNumber().ToString("D5"));
|
||||
|
||||
this.log = log;
|
||||
this.rootFolder = rootFolder;
|
||||
}
|
||||
|
||||
public TrackedFile CreateEmptyFile(string label = "")
|
||||
@@ -41,10 +49,15 @@ namespace FileUtils
|
||||
return result;
|
||||
}
|
||||
|
||||
public TrackedFile GenerateFile(ByteSize size, string label)
|
||||
public TrackedFile GenerateFile(ByteSize size, string label = "")
|
||||
{
|
||||
return GenerateFile(o => o.Random(size), label);
|
||||
}
|
||||
|
||||
public TrackedFile GenerateFile(Action<IGenerateOption> options, string label = "")
|
||||
{
|
||||
var sw = Stopwatch.Begin(log);
|
||||
var result = GenerateRandomFile(size, label);
|
||||
var result = RunGenerators(options, label);
|
||||
sw.End($"Generated file {result.Describe()}.");
|
||||
return result;
|
||||
}
|
||||
@@ -57,16 +70,27 @@ namespace FileUtils
|
||||
public void ScopedFiles(Action action)
|
||||
{
|
||||
PushFileSet();
|
||||
action();
|
||||
PopFileSet();
|
||||
try
|
||||
{
|
||||
action();
|
||||
}
|
||||
finally
|
||||
{
|
||||
PopFileSet();
|
||||
}
|
||||
}
|
||||
|
||||
public T ScopedFiles<T>(Func<T> action)
|
||||
{
|
||||
PushFileSet();
|
||||
var result = action();
|
||||
PopFileSet();
|
||||
return result;
|
||||
try
|
||||
{
|
||||
return action();
|
||||
}
|
||||
finally
|
||||
{
|
||||
PopFileSet();
|
||||
}
|
||||
}
|
||||
|
||||
private void PushFileSet()
|
||||
@@ -89,26 +113,35 @@ namespace FileUtils
|
||||
if (!Directory.GetFiles(folder).Any()) DeleteDirectory();
|
||||
}
|
||||
|
||||
private TrackedFile GenerateRandomFile(ByteSize size, string label)
|
||||
private TrackedFile RunGenerators(Action<IGenerateOption> options, string label)
|
||||
{
|
||||
var result = CreateEmptyFile(label);
|
||||
CheckSpaceAvailable(result, size);
|
||||
var generators = GetGenerators(options);
|
||||
CheckSpaceAvailable(result, generators.GetRequiredSpace());
|
||||
|
||||
GenerateFileBytes(result, size);
|
||||
using var stream = new FileStream(result.Filename, FileMode.Append);
|
||||
generators.Run(stream);
|
||||
return result;
|
||||
}
|
||||
|
||||
private void CheckSpaceAvailable(TrackedFile testFile, ByteSize size)
|
||||
private GeneratorCollection GetGenerators(Action<IGenerateOption> options)
|
||||
{
|
||||
var result = new GeneratorCollection();
|
||||
options(result);
|
||||
return result;
|
||||
}
|
||||
|
||||
private void CheckSpaceAvailable(TrackedFile testFile, long requiredSize)
|
||||
{
|
||||
var file = new FileInfo(testFile.Filename);
|
||||
var drive = new DriveInfo(file.Directory!.Root.FullName);
|
||||
|
||||
var spaceAvailable = drive.TotalFreeSpace;
|
||||
|
||||
if (spaceAvailable < size.SizeInBytes)
|
||||
if (spaceAvailable < requiredSize)
|
||||
{
|
||||
var msg = $"Not enough disk space. " +
|
||||
$"{Formatter.FormatByteSize(size.SizeInBytes)} required. " +
|
||||
$"{Formatter.FormatByteSize(requiredSize)} required. " +
|
||||
$"{Formatter.FormatByteSize(spaceAvailable)} available.";
|
||||
|
||||
log.Log(msg);
|
||||
@@ -116,34 +149,6 @@ namespace FileUtils
|
||||
}
|
||||
}
|
||||
|
||||
private void GenerateFileBytes(TrackedFile result, ByteSize size)
|
||||
{
|
||||
long bytesLeft = size.SizeInBytes;
|
||||
int chunkSize = ChunkSize;
|
||||
while (bytesLeft > 0)
|
||||
{
|
||||
try
|
||||
{
|
||||
var length = Math.Min(bytesLeft, chunkSize);
|
||||
AppendRandomBytesToFile(result, length);
|
||||
bytesLeft -= length;
|
||||
}
|
||||
catch
|
||||
{
|
||||
chunkSize = chunkSize / 2;
|
||||
if (chunkSize < 1024) throw;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private void AppendRandomBytesToFile(TrackedFile result, long length)
|
||||
{
|
||||
var bytes = new byte[length];
|
||||
random.NextBytes(bytes);
|
||||
using var stream = new FileStream(result.Filename, FileMode.Append);
|
||||
stream.Write(bytes, 0, bytes.Length);
|
||||
}
|
||||
|
||||
private void EnsureDirectory()
|
||||
{
|
||||
Directory.CreateDirectory(folder);
|
||||
|
||||
@@ -0,0 +1,141 @@
|
||||
using System.Text;
|
||||
using Utils;
|
||||
|
||||
namespace FileUtils
|
||||
{
|
||||
public class GeneratorCollection : IGenerateOption
|
||||
{
|
||||
private readonly List<IGenerator> generators = new List<IGenerator>();
|
||||
|
||||
public IGenerateOption ByteRepeat(byte[] bytes, ByteSize size)
|
||||
{
|
||||
var times = size.SizeInBytes / bytes.Length;
|
||||
generators.Add(new ByteRepeater(bytes, times));
|
||||
return this;
|
||||
}
|
||||
|
||||
public IGenerateOption ByteRepeat(byte[] bytes, int times)
|
||||
{
|
||||
generators.Add(new ByteRepeater(bytes, times));
|
||||
return this;
|
||||
}
|
||||
|
||||
public IGenerateOption Random(ByteSize size)
|
||||
{
|
||||
generators.Add(new RandomGenerator(size));
|
||||
return this;
|
||||
}
|
||||
|
||||
public IGenerateOption StringRepeat(string str, ByteSize size)
|
||||
{
|
||||
var times = size.SizeInBytes / str.Length;
|
||||
generators.Add(new StringRepeater(str, times));
|
||||
return this;
|
||||
}
|
||||
|
||||
public IGenerateOption StringRepeat(string str, int times)
|
||||
{
|
||||
generators.Add(new StringRepeater(str, times));
|
||||
return this;
|
||||
}
|
||||
|
||||
public void Run(FileStream file)
|
||||
{
|
||||
foreach (var generator in generators)
|
||||
{
|
||||
generator.Generate(file);
|
||||
}
|
||||
}
|
||||
|
||||
public long GetRequiredSpace()
|
||||
{
|
||||
return generators.Sum(g => g.GetRequiredSpace());
|
||||
}
|
||||
}
|
||||
|
||||
public interface IGenerator
|
||||
{
|
||||
void Generate(FileStream file);
|
||||
long GetRequiredSpace();
|
||||
}
|
||||
|
||||
public class ByteRepeater : IGenerator
|
||||
{
|
||||
private readonly byte[] bytes;
|
||||
private readonly long times;
|
||||
|
||||
public ByteRepeater(byte[] bytes, long times)
|
||||
{
|
||||
this.bytes = bytes;
|
||||
this.times = times;
|
||||
}
|
||||
|
||||
public void Generate(FileStream file)
|
||||
{
|
||||
for (var i = 0; i < times; i++)
|
||||
{
|
||||
file.Write(bytes, 0, bytes.Length);
|
||||
}
|
||||
}
|
||||
|
||||
public long GetRequiredSpace()
|
||||
{
|
||||
return bytes.Length * times;
|
||||
}
|
||||
}
|
||||
|
||||
public class StringRepeater : IGenerator
|
||||
{
|
||||
private readonly string str;
|
||||
private readonly long times;
|
||||
|
||||
public StringRepeater(string str, long times)
|
||||
{
|
||||
this.str = str;
|
||||
this.times = times;
|
||||
}
|
||||
|
||||
public void Generate(FileStream file)
|
||||
{
|
||||
using var writer = new StreamWriter(file);
|
||||
for (var i = 0; i < times; i++)
|
||||
{
|
||||
writer.Write(str);
|
||||
}
|
||||
}
|
||||
|
||||
public long GetRequiredSpace()
|
||||
{
|
||||
return Encoding.ASCII.GetBytes(str).Length * times;
|
||||
}
|
||||
}
|
||||
|
||||
public class RandomGenerator : IGenerator
|
||||
{
|
||||
private readonly Random random = new Random();
|
||||
private readonly ByteSize size;
|
||||
|
||||
public RandomGenerator(ByteSize size)
|
||||
{
|
||||
this.size = size;
|
||||
}
|
||||
|
||||
public void Generate(FileStream file)
|
||||
{
|
||||
var bytesLeft = size.SizeInBytes;
|
||||
while (bytesLeft > 0)
|
||||
{
|
||||
var size = Math.Min(bytesLeft, FileManager.ChunkSize);
|
||||
var bytes = new byte[size];
|
||||
random.NextBytes(bytes);
|
||||
file.Write(bytes, 0, bytes.Length);
|
||||
bytesLeft -= size;
|
||||
}
|
||||
}
|
||||
|
||||
public long GetRequiredSpace()
|
||||
{
|
||||
return size.SizeInBytes;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,13 +1,15 @@
|
||||
using KubernetesWorkflow;
|
||||
using Logging;
|
||||
using Logging;
|
||||
|
||||
namespace Core
|
||||
namespace KubernetesWorkflow
|
||||
{
|
||||
public interface IDownloadedLog
|
||||
{
|
||||
void IterateLines(Action<string> action);
|
||||
string ContainerName { get; }
|
||||
|
||||
void IterateLines(Action<string> action, params string[] thatContain);
|
||||
string[] GetLinesContaining(string expectedString);
|
||||
string[] FindLinesThatContain(params string[] tags);
|
||||
string GetFilepath();
|
||||
void DeleteFile();
|
||||
}
|
||||
|
||||
@@ -15,12 +17,15 @@ namespace Core
|
||||
{
|
||||
private readonly LogFile logFile;
|
||||
|
||||
internal DownloadedLog(WriteToFileLogHandler logHandler)
|
||||
internal DownloadedLog(WriteToFileLogHandler logHandler, string containerName)
|
||||
{
|
||||
logFile = logHandler.LogFile;
|
||||
ContainerName = containerName;
|
||||
}
|
||||
|
||||
public void IterateLines(Action<string> action)
|
||||
public string ContainerName { get; }
|
||||
|
||||
public void IterateLines(Action<string> action, params string[] thatContain)
|
||||
{
|
||||
using var file = File.OpenRead(logFile.FullFilename);
|
||||
using var streamReader = new StreamReader(file);
|
||||
@@ -28,7 +33,10 @@ namespace Core
|
||||
var line = streamReader.ReadLine();
|
||||
while (line != null)
|
||||
{
|
||||
action(line);
|
||||
if (thatContain.All(line.Contains))
|
||||
{
|
||||
action(line);
|
||||
}
|
||||
line = streamReader.ReadLine();
|
||||
}
|
||||
}
|
||||
@@ -72,6 +80,11 @@ namespace Core
|
||||
return result.ToArray();
|
||||
}
|
||||
|
||||
public string GetFilepath()
|
||||
{
|
||||
return logFile.FullFilename;
|
||||
}
|
||||
|
||||
public void DeleteFile()
|
||||
{
|
||||
File.Delete(logFile.FullFilename);
|
||||
@@ -1,4 +1,5 @@
|
||||
using Newtonsoft.Json.Linq;
|
||||
using Newtonsoft.Json;
|
||||
using Newtonsoft.Json.Linq;
|
||||
|
||||
namespace KubernetesWorkflow.Recipe
|
||||
{
|
||||
@@ -21,14 +22,13 @@ namespace KubernetesWorkflow.Recipe
|
||||
var typeName = GetTypeName(typeof(T));
|
||||
var userData = Additionals.SingleOrDefault(a => a.Type == typeName);
|
||||
if (userData == null) return default;
|
||||
var jobject = (JObject)userData.UserData;
|
||||
return jobject.ToObject<T>();
|
||||
return JsonConvert.DeserializeObject<T>(userData.UserData);
|
||||
}
|
||||
|
||||
private static Additional ConvertToAdditional(object userData)
|
||||
{
|
||||
var typeName = GetTypeName(userData.GetType());
|
||||
return new Additional(typeName, userData);
|
||||
return new Additional(typeName, JsonConvert.SerializeObject(userData));
|
||||
}
|
||||
|
||||
private static string GetTypeName(Type type)
|
||||
@@ -41,13 +41,13 @@ namespace KubernetesWorkflow.Recipe
|
||||
|
||||
public class Additional
|
||||
{
|
||||
public Additional(string type, object userData)
|
||||
public Additional(string type, string userData)
|
||||
{
|
||||
Type = type;
|
||||
UserData = userData;
|
||||
}
|
||||
|
||||
public string Type { get; }
|
||||
public object UserData { get; }
|
||||
public string UserData { get; }
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2,8 +2,9 @@
|
||||
{
|
||||
public class ContainerRecipe
|
||||
{
|
||||
public ContainerRecipe(int number, string? nameOverride, string image, ContainerResources resources, SchedulingAffinity schedulingAffinity, CommandOverride commandOverride, bool setCriticalPriority, Port[] exposedPorts, Port[] internalPorts, EnvVar[] envVars, PodLabels podLabels, PodAnnotations podAnnotations, VolumeMount[] volumes, ContainerAdditionals additionals)
|
||||
public ContainerRecipe(DateTime recipeCreatedUtc, int number, string? nameOverride, string image, ContainerResources resources, SchedulingAffinity schedulingAffinity, CommandOverride commandOverride, bool setCriticalPriority, Port[] exposedPorts, Port[] internalPorts, EnvVar[] envVars, PodLabels podLabels, PodAnnotations podAnnotations, VolumeMount[] volumes, ContainerAdditionals additionals)
|
||||
{
|
||||
RecipeCreatedUtc = recipeCreatedUtc;
|
||||
Number = number;
|
||||
NameOverride = nameOverride;
|
||||
Image = image;
|
||||
@@ -31,6 +32,7 @@
|
||||
if (exposedPorts.Any(p => string.IsNullOrEmpty(p.Tag))) throw new Exception("Port tags are required for all exposed ports.");
|
||||
}
|
||||
|
||||
public DateTime RecipeCreatedUtc { get; }
|
||||
public string Name { get; }
|
||||
public int Number { get; }
|
||||
public string? NameOverride { get; }
|
||||
|
||||
@@ -25,7 +25,7 @@ namespace KubernetesWorkflow.Recipe
|
||||
|
||||
Initialize(config);
|
||||
|
||||
var recipe = new ContainerRecipe(containerNumber, config.NameOverride, Image, resources, schedulingAffinity, commandOverride, setCriticalPriority,
|
||||
var recipe = new ContainerRecipe(DateTime.UtcNow, containerNumber, config.NameOverride, Image, resources, schedulingAffinity, commandOverride, setCriticalPriority,
|
||||
exposedPorts.ToArray(),
|
||||
internalPorts.ToArray(),
|
||||
envVars.ToArray(),
|
||||
|
||||
@@ -16,6 +16,7 @@ namespace KubernetesWorkflow
|
||||
CrashWatcher CreateCrashWatcher(RunningContainer container);
|
||||
void Stop(RunningPod pod, bool waitTillStopped);
|
||||
void DownloadContainerLog(RunningContainer container, ILogHandler logHandler, int? tailLines = null, bool? previous = null);
|
||||
IDownloadedLog DownloadContainerLog(RunningContainer container, int? tailLines = null, bool? previous = null);
|
||||
string ExecuteCommand(RunningContainer container, string command, params string[] args);
|
||||
void DeleteNamespace(bool wait);
|
||||
void DeleteNamespacesStartingWith(string namespacePrefix, bool wait);
|
||||
@@ -60,7 +61,7 @@ namespace KubernetesWorkflow
|
||||
var startResult = controller.BringOnline(recipes, location);
|
||||
var containers = CreateContainers(startResult, recipes, startupConfig);
|
||||
|
||||
var rc = new RunningPod(startupConfig, startResult, containers);
|
||||
var rc = new RunningPod(Guid.NewGuid().ToString(), startupConfig, startResult, containers);
|
||||
cluster.Configuration.Hooks.OnContainersStarted(rc);
|
||||
|
||||
if (startResult.ExternalService != null)
|
||||
@@ -99,11 +100,19 @@ namespace KubernetesWorkflow
|
||||
|
||||
public void Stop(RunningPod runningPod, bool waitTillStopped)
|
||||
{
|
||||
if (runningPod.IsStopped) return;
|
||||
foreach (var c in runningPod.Containers)
|
||||
{
|
||||
c.StopLog = DownloadContainerLog(c);
|
||||
}
|
||||
runningPod.IsStopped = true;
|
||||
|
||||
K8s(controller =>
|
||||
{
|
||||
controller.Stop(runningPod.StartResult, waitTillStopped);
|
||||
cluster.Configuration.Hooks.OnContainersStopped(runningPod);
|
||||
});
|
||||
|
||||
cluster.Configuration.Hooks.OnContainersStopped(runningPod);
|
||||
}
|
||||
|
||||
public void DownloadContainerLog(RunningContainer container, ILogHandler logHandler, int? tailLines = null, bool? previous = null)
|
||||
@@ -114,6 +123,20 @@ namespace KubernetesWorkflow
|
||||
});
|
||||
}
|
||||
|
||||
public IDownloadedLog DownloadContainerLog(RunningContainer container, int? tailLines = null, bool? previous = null)
|
||||
{
|
||||
var msg = $"Downloading container log for '{container.Name}'";
|
||||
log.Log(msg);
|
||||
var logHandler = new WriteToFileLogHandler(log, msg);
|
||||
|
||||
K8s(controller =>
|
||||
{
|
||||
controller.DownloadPodLog(container, logHandler, tailLines, previous);
|
||||
});
|
||||
|
||||
return new DownloadedLog(logHandler, container.Name);
|
||||
}
|
||||
|
||||
public string ExecuteCommand(RunningContainer container, string command, params string[] args)
|
||||
{
|
||||
return K8s(controller =>
|
||||
@@ -147,7 +170,7 @@ namespace KubernetesWorkflow
|
||||
var addresses = CreateContainerAddresses(startResult, r);
|
||||
log.Debug($"{r}={name} -> container addresses: {string.Join(Environment.NewLine, addresses.Select(a => a.ToString()))}");
|
||||
|
||||
return new RunningContainer(name, r, addresses);
|
||||
return new RunningContainer(Guid.NewGuid().ToString(), name, r, addresses);
|
||||
|
||||
}).ToArray();
|
||||
}
|
||||
|
||||
@@ -7,16 +7,19 @@ namespace KubernetesWorkflow.Types
|
||||
{
|
||||
public class RunningContainer
|
||||
{
|
||||
public RunningContainer(string name, ContainerRecipe recipe, ContainerAddress[] addresses)
|
||||
public RunningContainer(string id, string name, ContainerRecipe recipe, ContainerAddress[] addresses)
|
||||
{
|
||||
Id = id;
|
||||
Name = name;
|
||||
Recipe = recipe;
|
||||
Addresses = addresses;
|
||||
}
|
||||
|
||||
public string Id { get; }
|
||||
public string Name { get; }
|
||||
public ContainerRecipe Recipe { get; }
|
||||
public ContainerAddress[] Addresses { get; }
|
||||
public IDownloadedLog? StopLog { get; internal set; }
|
||||
|
||||
[JsonIgnore]
|
||||
public RunningPod RunningPod { get; internal set; } = null!;
|
||||
@@ -50,5 +53,21 @@ namespace KubernetesWorkflow.Types
|
||||
}
|
||||
throw new Exception("Running location not known.");
|
||||
}
|
||||
|
||||
public override string ToString()
|
||||
{
|
||||
return Name;
|
||||
}
|
||||
|
||||
public override bool Equals(object? obj)
|
||||
{
|
||||
return obj is RunningContainer container &&
|
||||
Id == container.Id;
|
||||
}
|
||||
|
||||
public override int GetHashCode()
|
||||
{
|
||||
return HashCode.Combine(Id);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -4,8 +4,9 @@ namespace KubernetesWorkflow.Types
|
||||
{
|
||||
public class RunningPod
|
||||
{
|
||||
public RunningPod(StartupConfig startupConfig, StartResult startResult, RunningContainer[] containers)
|
||||
public RunningPod(string id, StartupConfig startupConfig, StartResult startResult, RunningContainer[] containers)
|
||||
{
|
||||
Id = id;
|
||||
StartupConfig = startupConfig;
|
||||
StartResult = startResult;
|
||||
Containers = containers;
|
||||
@@ -13,6 +14,7 @@ namespace KubernetesWorkflow.Types
|
||||
foreach (var c in containers) c.RunningPod = this;
|
||||
}
|
||||
|
||||
public string Id { get; }
|
||||
public StartupConfig StartupConfig { get; }
|
||||
public StartResult StartResult { get; }
|
||||
public RunningContainer[] Containers { get; }
|
||||
@@ -23,10 +25,30 @@ namespace KubernetesWorkflow.Types
|
||||
get { return $"'{string.Join("&", Containers.Select(c => c.Name).ToArray())}'"; }
|
||||
}
|
||||
|
||||
[JsonIgnore]
|
||||
public bool IsStopped { get; internal set; }
|
||||
|
||||
public string Describe()
|
||||
{
|
||||
return string.Join(",", Containers.Select(c => c.Name));
|
||||
}
|
||||
|
||||
public override bool Equals(object? obj)
|
||||
{
|
||||
return obj is RunningPod pod &&
|
||||
Id == pod.Id;
|
||||
}
|
||||
|
||||
public override int GetHashCode()
|
||||
{
|
||||
return HashCode.Combine(Id);
|
||||
}
|
||||
|
||||
public override string ToString()
|
||||
{
|
||||
if (IsStopped) return Name + " (*)";
|
||||
return Name;
|
||||
}
|
||||
}
|
||||
|
||||
public static class RunningContainersExtensions
|
||||
|
||||
@@ -77,7 +77,7 @@ namespace Logging
|
||||
return new LogFile($"{GetFullName()}_{GetSubfileNumber()}", ext);
|
||||
}
|
||||
|
||||
private string ApplyReplacements(string str)
|
||||
protected string ApplyReplacements(string str)
|
||||
{
|
||||
if (IsDebug) return str;
|
||||
foreach (var replacement in replacements)
|
||||
|
||||
@@ -9,7 +9,7 @@
|
||||
|
||||
public override void Log(string message)
|
||||
{
|
||||
Console.WriteLine(message);
|
||||
Console.WriteLine(ApplyReplacements(message));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -117,6 +117,7 @@ namespace NethereumWorkflow
|
||||
}
|
||||
|
||||
return new BlockInterval(
|
||||
timeRange: timeRange,
|
||||
from: fromBlock.Value,
|
||||
to: toBlock.Value
|
||||
);
|
||||
|
||||
@@ -0,0 +1,61 @@
|
||||
namespace OverwatchTranscript
|
||||
{
|
||||
public class ActionQueue
|
||||
{
|
||||
// Using ConcurrentQueue<> here would make this process slower.
|
||||
private readonly object queueLock = new object();
|
||||
private readonly AutoResetEvent signal = new AutoResetEvent(false);
|
||||
private List<Action> queue = new List<Action>();
|
||||
private Task queueWorker = null!;
|
||||
private bool stopping = false;
|
||||
|
||||
public void Start()
|
||||
{
|
||||
queueWorker = Task.Run(QueueWorker);
|
||||
}
|
||||
|
||||
public int Count { get; private set; }
|
||||
|
||||
public void StopAndJoin()
|
||||
{
|
||||
stopping = true;
|
||||
queueWorker.Wait();
|
||||
if (queue.Count > 0) throw new Exception("not all acions handled");
|
||||
queueWorker.Dispose();
|
||||
}
|
||||
|
||||
public void Add(Action action)
|
||||
{
|
||||
if (stopping) throw new Exception("queue stopping");
|
||||
|
||||
lock (queueLock)
|
||||
{
|
||||
queue.Add(action);
|
||||
Count = queue.Count;
|
||||
}
|
||||
signal.Set();
|
||||
}
|
||||
|
||||
private void QueueWorker()
|
||||
{
|
||||
while (true)
|
||||
{
|
||||
signal.WaitOne(10);
|
||||
|
||||
List<Action> work = null!;
|
||||
lock (queueLock)
|
||||
{
|
||||
work = queue;
|
||||
queue = new List<Action>();
|
||||
Count = 0;
|
||||
}
|
||||
if (stopping && !work.Any()) return;
|
||||
|
||||
foreach (var action in work)
|
||||
{
|
||||
action();
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,110 @@
|
||||
using Logging;
|
||||
using System.Collections.Concurrent;
|
||||
|
||||
namespace OverwatchTranscript
|
||||
{
|
||||
public class BucketSet
|
||||
{
|
||||
private const int numberOfActiveBuckets = 10;
|
||||
private readonly ILog log;
|
||||
private readonly string workingDir;
|
||||
private readonly object _bucketLock = new object();
|
||||
private readonly List<EventBucketWriter> fullBuckets = new List<EventBucketWriter>();
|
||||
private readonly List<EventBucketWriter> activeBuckets = new List<EventBucketWriter>();
|
||||
private readonly ActionQueue queue = new ActionQueue();
|
||||
private int activeBucketIndex = 0;
|
||||
private bool closed = false;
|
||||
private string internalErrors = string.Empty;
|
||||
|
||||
public BucketSet(ILog log, string workingDir)
|
||||
{
|
||||
this.log = log;
|
||||
this.workingDir = workingDir;
|
||||
|
||||
for (var i = 0; i < numberOfActiveBuckets;i++)
|
||||
{
|
||||
AddNewBucket();
|
||||
}
|
||||
|
||||
queue.Start();
|
||||
}
|
||||
|
||||
public void Add(DateTime utc, object payload)
|
||||
{
|
||||
if (closed) throw new Exception("Buckets already closed!");
|
||||
queue.Add(() => AddInternal(utc, payload));
|
||||
|
||||
if (queue.Count > 1000)
|
||||
{
|
||||
Thread.Sleep(1);
|
||||
}
|
||||
}
|
||||
|
||||
public IFinalizedBucket[] FinalizeBuckets()
|
||||
{
|
||||
closed = true;
|
||||
queue.StopAndJoin();
|
||||
|
||||
if (IsEmpty()) throw new Exception("No entries have been added.");
|
||||
if (!string.IsNullOrEmpty(internalErrors)) throw new Exception(internalErrors);
|
||||
|
||||
var buckets = fullBuckets.Concat(activeBuckets).ToArray();
|
||||
log.Debug($"Finalizing {buckets.Length} buckets...");
|
||||
|
||||
var finalized = new ConcurrentBag<IFinalizedBucket>();
|
||||
var tasks = Parallel.ForEach(buckets, b => finalized.Add(b.FinalizeBucket()));
|
||||
if (!tasks.IsCompleted) throw new Exception("Failed to finalize buckets: " + tasks);
|
||||
|
||||
return finalized.ToArray();
|
||||
}
|
||||
|
||||
private bool IsEmpty()
|
||||
{
|
||||
return fullBuckets.All(b => b.Count == 0) && activeBuckets.All(b => b.Count == 0);
|
||||
}
|
||||
|
||||
private void AddInternal(DateTime utc, object payload)
|
||||
{
|
||||
try
|
||||
{
|
||||
lock (_bucketLock)
|
||||
{
|
||||
var current = activeBuckets[activeBucketIndex];
|
||||
current.Add(utc, payload);
|
||||
activeBucketIndex = (activeBucketIndex + 1) % numberOfActiveBuckets;
|
||||
|
||||
if (current.IsFull)
|
||||
{
|
||||
log.Debug("Bucket is full. New bucket...");
|
||||
fullBuckets.Add(current);
|
||||
activeBuckets.Remove(current);
|
||||
AddNewBucket();
|
||||
}
|
||||
}
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
internalErrors += ex.ToString();
|
||||
log.Error(ex.ToString());
|
||||
}
|
||||
}
|
||||
|
||||
private static int bucketSizeIndex = 0;
|
||||
private static int[] bucketSizes = new[]
|
||||
{
|
||||
10000,
|
||||
15000,
|
||||
20000,
|
||||
};
|
||||
|
||||
private void AddNewBucket()
|
||||
{
|
||||
lock (_bucketLock)
|
||||
{
|
||||
var size = bucketSizes[bucketSizeIndex];
|
||||
bucketSizeIndex = (bucketSizeIndex + 1) % bucketSizes.Length;
|
||||
activeBuckets.Add(new EventBucketWriter(log, Path.Combine(workingDir, Guid.NewGuid().ToString()), size));
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,143 @@
|
||||
using Logging;
|
||||
using Newtonsoft.Json;
|
||||
using System.Collections.Concurrent;
|
||||
|
||||
namespace OverwatchTranscript
|
||||
{
|
||||
public interface IFinalizedBucket
|
||||
{
|
||||
bool IsEmpty { get; }
|
||||
DateTime? SeeTopUtc();
|
||||
BucketTop? TakeTop();
|
||||
}
|
||||
|
||||
public class BucketTop
|
||||
{
|
||||
public BucketTop(DateTime utc, OverwatchEvent[] events)
|
||||
{
|
||||
Utc = utc;
|
||||
Events = events;
|
||||
}
|
||||
|
||||
public DateTime Utc { get; }
|
||||
public OverwatchEvent[] Events { get; }
|
||||
}
|
||||
|
||||
public class EventBucketReader : IFinalizedBucket
|
||||
{
|
||||
private readonly string bucketFile;
|
||||
private readonly ConcurrentQueue<BucketTop> topQueue = new ConcurrentQueue<BucketTop>();
|
||||
private readonly AutoResetEvent itemDequeued = new AutoResetEvent(false);
|
||||
private bool stopping;
|
||||
|
||||
public EventBucketReader(ILog log, string bucketFile)
|
||||
{
|
||||
this.bucketFile = bucketFile;
|
||||
if (!File.Exists(bucketFile)) throw new Exception("Doesn't exist: " + bucketFile);
|
||||
|
||||
log.Debug("Read Bucket open: " + bucketFile);
|
||||
|
||||
Task.Run(ReadBucket);
|
||||
}
|
||||
|
||||
public bool IsEmpty { get; private set; }
|
||||
|
||||
public DateTime? SeeTopUtc()
|
||||
{
|
||||
if (IsEmpty) return null;
|
||||
while (true)
|
||||
{
|
||||
UpdateIsEmpty();
|
||||
if (IsEmpty) return null;
|
||||
if (topQueue.TryPeek(out BucketTop? top))
|
||||
{
|
||||
return top.Utc;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
public BucketTop? TakeTop()
|
||||
{
|
||||
if (IsEmpty) return null;
|
||||
|
||||
while (true)
|
||||
{
|
||||
UpdateIsEmpty();
|
||||
if (IsEmpty) return null;
|
||||
if (topQueue.TryDequeue(out BucketTop? top))
|
||||
{
|
||||
itemDequeued.Set();
|
||||
return top;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private void ReadBucket()
|
||||
{
|
||||
using var file = File.OpenRead(bucketFile);
|
||||
using var reader = new StreamReader(file);
|
||||
|
||||
while (true)
|
||||
{
|
||||
while (topQueue.Count < 5)
|
||||
{
|
||||
var top = CreateNewTop(reader);
|
||||
if (top != null)
|
||||
{
|
||||
topQueue.Enqueue(top);
|
||||
}
|
||||
else
|
||||
{
|
||||
stopping = true;
|
||||
return;
|
||||
}
|
||||
}
|
||||
|
||||
itemDequeued.Reset();
|
||||
itemDequeued.WaitOne();
|
||||
}
|
||||
}
|
||||
|
||||
private void UpdateIsEmpty()
|
||||
{
|
||||
var empty = stopping && topQueue.IsEmpty;
|
||||
if (!IsEmpty && empty)
|
||||
{
|
||||
File.Delete(bucketFile);
|
||||
IsEmpty = true;
|
||||
}
|
||||
}
|
||||
|
||||
private EventBucketEntry? nextEntry = null;
|
||||
private BucketTop? CreateNewTop(StreamReader reader)
|
||||
{
|
||||
if (nextEntry == null)
|
||||
{
|
||||
nextEntry = ReadEntry(reader);
|
||||
if (nextEntry == null) return null;
|
||||
}
|
||||
|
||||
var topEntry = nextEntry;
|
||||
var entries = new List<EventBucketEntry>
|
||||
{
|
||||
topEntry
|
||||
};
|
||||
|
||||
nextEntry = ReadEntry(reader);
|
||||
while (nextEntry != null && nextEntry.Utc == topEntry.Utc)
|
||||
{
|
||||
entries.Add(nextEntry);
|
||||
nextEntry = ReadEntry(reader);
|
||||
}
|
||||
|
||||
return new BucketTop(topEntry.Utc, entries.Select(e => e.Event).ToArray());
|
||||
}
|
||||
|
||||
private EventBucketEntry? ReadEntry(StreamReader reader)
|
||||
{
|
||||
var line = reader.ReadLine();
|
||||
if (string.IsNullOrEmpty(line)) return null;
|
||||
return JsonConvert.DeserializeObject<EventBucketEntry>(line);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,113 @@
|
||||
using Logging;
|
||||
using Newtonsoft.Json;
|
||||
|
||||
namespace OverwatchTranscript
|
||||
{
|
||||
public class EventBucketWriter
|
||||
{
|
||||
private const int MaxBuffer = 1000;
|
||||
|
||||
private readonly object _lock = new object();
|
||||
private bool closed = false;
|
||||
private readonly ILog log;
|
||||
private readonly string bucketFile;
|
||||
private readonly int maxCount;
|
||||
private readonly List<EventBucketEntry> buffer = new List<EventBucketEntry>();
|
||||
|
||||
public EventBucketWriter(ILog log, string bucketFile, int maxCount)
|
||||
{
|
||||
this.log = log;
|
||||
this.bucketFile = bucketFile;
|
||||
this.maxCount = maxCount;
|
||||
if (File.Exists(bucketFile)) throw new Exception("Already exists");
|
||||
|
||||
log.Debug("Write Bucket open: " + bucketFile);
|
||||
}
|
||||
|
||||
public int Count { get; private set; }
|
||||
public bool IsFull { get; private set; }
|
||||
|
||||
public void Add(DateTime utc, object payload)
|
||||
{
|
||||
lock (_lock)
|
||||
{
|
||||
if (closed) throw new Exception("Already closed");
|
||||
AddToBuffer(utc, payload);
|
||||
BufferToFile(emptyBuffer: false);
|
||||
}
|
||||
}
|
||||
|
||||
public IFinalizedBucket FinalizeBucket()
|
||||
{
|
||||
lock (_lock)
|
||||
{
|
||||
closed = true;
|
||||
BufferToFile(emptyBuffer: true);
|
||||
SortFileByTimestamps();
|
||||
}
|
||||
log.Debug($"Finalized bucket with {Count} entries");
|
||||
return new EventBucketReader(log, bucketFile);
|
||||
}
|
||||
|
||||
public override string ToString()
|
||||
{
|
||||
return $"EventBucket: " + Count;
|
||||
}
|
||||
|
||||
private void AddToBuffer(DateTime utc, object payload)
|
||||
{
|
||||
var typeName = payload.GetType().FullName;
|
||||
if (string.IsNullOrEmpty(typeName)) throw new Exception("Empty typename for payload");
|
||||
if (utc == default) throw new Exception("DateTimeUtc not set");
|
||||
|
||||
var entry = new EventBucketEntry
|
||||
{
|
||||
Utc = utc,
|
||||
Event = new OverwatchEvent
|
||||
{
|
||||
Type = typeName,
|
||||
Payload = Json.Serialize(payload)
|
||||
}
|
||||
};
|
||||
|
||||
Count++;
|
||||
IsFull = Count > maxCount;
|
||||
|
||||
buffer.Add(entry);
|
||||
}
|
||||
|
||||
private void BufferToFile(bool emptyBuffer)
|
||||
{
|
||||
if (emptyBuffer || buffer.Count > MaxBuffer)
|
||||
{
|
||||
using var file = File.Open(bucketFile, FileMode.Append);
|
||||
using var writer = new StreamWriter(file);
|
||||
foreach (var entry in buffer)
|
||||
{
|
||||
writer.WriteLine(Json.Serialize(entry));
|
||||
}
|
||||
log.Debug($"Bucket wrote {buffer.Count} entries to file.");
|
||||
buffer.Clear();
|
||||
}
|
||||
}
|
||||
|
||||
private void SortFileByTimestamps()
|
||||
{
|
||||
var lines = File.ReadAllLines(bucketFile);
|
||||
var entries = lines.Select(Json.Deserialize<EventBucketEntry>)
|
||||
.Cast<EventBucketEntry>()
|
||||
.OrderBy(e => e.Utc)
|
||||
.ToArray();
|
||||
|
||||
File.Delete(bucketFile);
|
||||
File.WriteAllLines(bucketFile, entries.Select(e => Json.Serialize(e)));
|
||||
}
|
||||
}
|
||||
|
||||
[Serializable]
|
||||
public class EventBucketEntry
|
||||
{
|
||||
public DateTime Utc { get; set; }
|
||||
public OverwatchEvent Event { get; set; } = new();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,27 @@
|
||||
using Newtonsoft.Json;
|
||||
using System.Globalization;
|
||||
|
||||
namespace OverwatchTranscript
|
||||
{
|
||||
public static class Json
|
||||
{
|
||||
private static JsonSerializerSettings settings = new JsonSerializerSettings
|
||||
{
|
||||
Formatting = Formatting.None,
|
||||
NullValueHandling = NullValueHandling.Ignore,
|
||||
Culture = CultureInfo.InvariantCulture,
|
||||
DateFormatHandling = DateFormatHandling.IsoDateFormat,
|
||||
FloatFormatHandling = FloatFormatHandling.Symbol
|
||||
};
|
||||
|
||||
public static string Serialize(object obj, Formatting formatting = Formatting.None)
|
||||
{
|
||||
return JsonConvert.SerializeObject(obj, formatting, settings);
|
||||
}
|
||||
|
||||
public static T Deserialize<T>(string json)
|
||||
{
|
||||
return JsonConvert.DeserializeObject<T>(json)!;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,56 @@
|
||||
namespace OverwatchTranscript
|
||||
{
|
||||
[Serializable]
|
||||
public class OverwatchTranscript
|
||||
{
|
||||
public OverwatchHeader Header { get; set; } = new();
|
||||
public OverwatchMomentReference[] MomentReferences { get; set; } = Array.Empty<OverwatchMomentReference>();
|
||||
}
|
||||
|
||||
[Serializable]
|
||||
public class OverwatchMomentReference
|
||||
{
|
||||
public string MomentsFile { get; set; } = string.Empty;
|
||||
public int NumberOfMoments { get; set; }
|
||||
public int NumberOfEvents { get; set; }
|
||||
public DateTime EarliestUtc { get; set; }
|
||||
public DateTime LatestUtc { get; set; }
|
||||
}
|
||||
|
||||
[Serializable]
|
||||
public class OverwatchHeader
|
||||
{
|
||||
public OverwatchCommonHeader Common { get; set; } = new();
|
||||
public OverwatchHeaderEntry[] Entries { get; set; } = Array.Empty<OverwatchHeaderEntry>();
|
||||
}
|
||||
|
||||
[Serializable]
|
||||
public class OverwatchCommonHeader
|
||||
{
|
||||
public long NumberOfMoments { get; set; }
|
||||
public long NumberOfEvents { get; set; }
|
||||
public DateTime EarliestUtc { get; set; }
|
||||
public DateTime LatestUtc { get; set; }
|
||||
}
|
||||
|
||||
[Serializable]
|
||||
public class OverwatchHeaderEntry
|
||||
{
|
||||
public string Key { get; set; } = string.Empty;
|
||||
public string Value { get; set; } = string.Empty;
|
||||
}
|
||||
|
||||
[Serializable]
|
||||
public class OverwatchMoment
|
||||
{
|
||||
public DateTime Utc { get; set; }
|
||||
public OverwatchEvent[] Events { get; set; } = Array.Empty<OverwatchEvent>();
|
||||
}
|
||||
|
||||
[Serializable]
|
||||
public class OverwatchEvent
|
||||
{
|
||||
public string Type { get; set; } = string.Empty;
|
||||
public string Payload { get; set; } = string.Empty;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,104 @@
|
||||
using Newtonsoft.Json;
|
||||
|
||||
namespace OverwatchTranscript
|
||||
{
|
||||
public class MomentReader
|
||||
{
|
||||
private readonly OverwatchTranscript model;
|
||||
private readonly string workingDir;
|
||||
private int referenceIndex = 0;
|
||||
private int momentsRead = 0;
|
||||
private OpenReference currentRef;
|
||||
|
||||
public MomentReader(OverwatchTranscript model, string workingDir)
|
||||
{
|
||||
this.model = model;
|
||||
this.workingDir = workingDir;
|
||||
|
||||
currentRef = CreateOpenReference();
|
||||
}
|
||||
|
||||
public OverwatchMoment? Next()
|
||||
{
|
||||
if (referenceIndex >= model.MomentReferences.Length) return null;
|
||||
|
||||
var moment = currentRef.ReadNext();
|
||||
if (moment == null)
|
||||
{
|
||||
Close();
|
||||
|
||||
// This reference file ran out.
|
||||
// The number of moments read should match exactly the number of moments
|
||||
// describe in the reference. If not, error:
|
||||
var expected = model.MomentReferences[referenceIndex].NumberOfMoments;
|
||||
if (momentsRead != expected)
|
||||
{
|
||||
throw new Exception("Number of moments read from referenced file does not match number of moments value in model. " +
|
||||
$"Reads: { momentsRead} - model.MomentReferences[{referenceIndex}].NumberOfMoment: {expected}");
|
||||
}
|
||||
|
||||
referenceIndex++;
|
||||
if (referenceIndex < model.MomentReferences.Length)
|
||||
{
|
||||
// Proceed to next reference file.
|
||||
currentRef = CreateOpenReference();
|
||||
momentsRead = 0;
|
||||
return Next();
|
||||
}
|
||||
else
|
||||
{
|
||||
// That was the last one.
|
||||
return null;
|
||||
}
|
||||
}
|
||||
else
|
||||
{
|
||||
momentsRead++;
|
||||
return moment;
|
||||
}
|
||||
}
|
||||
|
||||
public void Close()
|
||||
{
|
||||
if (currentRef != null)
|
||||
{
|
||||
currentRef.Close();
|
||||
currentRef = null!;
|
||||
}
|
||||
}
|
||||
|
||||
private OpenReference CreateOpenReference()
|
||||
{
|
||||
var filepath = Path.Combine(workingDir, model.MomentReferences[referenceIndex].MomentsFile);
|
||||
return new OpenReference(filepath);
|
||||
}
|
||||
|
||||
private class OpenReference
|
||||
{
|
||||
private readonly FileStream file;
|
||||
private readonly StreamReader reader;
|
||||
|
||||
public OpenReference(string filePath)
|
||||
{
|
||||
file = File.OpenRead(filePath);
|
||||
reader = new StreamReader(file);
|
||||
}
|
||||
|
||||
public OverwatchMoment? ReadNext()
|
||||
{
|
||||
var line = reader.ReadLine();
|
||||
if (string.IsNullOrEmpty(line)) return null;
|
||||
return JsonConvert.DeserializeObject<OverwatchMoment>(line);
|
||||
}
|
||||
|
||||
public void Close()
|
||||
{
|
||||
reader.Close();
|
||||
file.Close();
|
||||
|
||||
reader.Dispose();
|
||||
file.Dispose();
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,151 @@
|
||||
using Logging;
|
||||
using Newtonsoft.Json;
|
||||
|
||||
namespace OverwatchTranscript
|
||||
{
|
||||
public class MomentReferenceBuilder
|
||||
{
|
||||
private const int MaxMomentsPerReference = 10000;
|
||||
private readonly ILog log;
|
||||
private readonly string workingDir;
|
||||
|
||||
public MomentReferenceBuilder(ILog log, string workingDir)
|
||||
{
|
||||
this.log = log;
|
||||
this.workingDir = workingDir;
|
||||
}
|
||||
|
||||
public OverwatchMomentReference[] Build(IFinalizedBucket[] finalizedBuckets)
|
||||
{
|
||||
var result = new List<OverwatchMomentReference>();
|
||||
var currentBuilder = new Builder(log, workingDir);
|
||||
|
||||
var buckets = finalizedBuckets.ToList();
|
||||
log.Debug($"Building references for {buckets.Count} buckets.");
|
||||
while (buckets.Any())
|
||||
{
|
||||
buckets.RemoveAll(b => b.IsEmpty);
|
||||
if (!buckets.Any()) break;
|
||||
|
||||
var earliestUtc = GetEarliestUtc(buckets);
|
||||
if (earliestUtc == null) continue;
|
||||
|
||||
var tops = CollectAllTopsForUtc(earliestUtc.Value, buckets);
|
||||
var moment = ConvertTopsToMoment(tops);
|
||||
currentBuilder.Add(moment);
|
||||
if (currentBuilder.NumberOfMoments == MaxMomentsPerReference)
|
||||
{
|
||||
result.Add(currentBuilder.Build());
|
||||
currentBuilder = new Builder(log, workingDir);
|
||||
}
|
||||
}
|
||||
|
||||
if (currentBuilder.NumberOfMoments > 0)
|
||||
{
|
||||
result.Add(currentBuilder.Build());
|
||||
}
|
||||
|
||||
return result.ToArray();
|
||||
}
|
||||
|
||||
private OverwatchMoment ConvertTopsToMoment(List<BucketTop> tops)
|
||||
{
|
||||
var discintUtc = tops.Select(e => e.Utc).Distinct().ToArray();
|
||||
if (discintUtc.Length != 1) throw new Exception("UTC mixing in moment construction.");
|
||||
|
||||
return new OverwatchMoment
|
||||
{
|
||||
Utc = tops[0].Utc,
|
||||
Events = tops.SelectMany(e => e.Events).ToArray()
|
||||
};
|
||||
}
|
||||
|
||||
private List<BucketTop> CollectAllTopsForUtc(DateTime earliestUtc, List<IFinalizedBucket> buckets)
|
||||
{
|
||||
var result = new List<BucketTop>();
|
||||
|
||||
foreach (var bucket in buckets)
|
||||
{
|
||||
if (bucket.IsEmpty) continue;
|
||||
|
||||
var utc = bucket.SeeTopUtc();
|
||||
if (utc == null) continue;
|
||||
|
||||
if (utc.Value == earliestUtc)
|
||||
{
|
||||
var top = bucket.TakeTop();
|
||||
if (top == null) throw new Exception("top was null after top utc was not");
|
||||
result.Add(top);
|
||||
}
|
||||
}
|
||||
|
||||
return result;
|
||||
}
|
||||
|
||||
private DateTime? GetEarliestUtc(List<IFinalizedBucket> buckets)
|
||||
{
|
||||
var earliest = DateTime.MaxValue;
|
||||
foreach (var bucket in buckets)
|
||||
{
|
||||
var utc = bucket.SeeTopUtc();
|
||||
if (utc == null) return null;
|
||||
|
||||
if (utc.Value < earliest) earliest = utc.Value;
|
||||
}
|
||||
return earliest;
|
||||
}
|
||||
|
||||
public class Builder
|
||||
{
|
||||
private readonly ILog log;
|
||||
private readonly string workingDir;
|
||||
private OverwatchMomentReference reference;
|
||||
private readonly ActionQueue queue = new ActionQueue();
|
||||
|
||||
public Builder(ILog log, string workingDir)
|
||||
{
|
||||
reference = new OverwatchMomentReference
|
||||
{
|
||||
MomentsFile = Guid.NewGuid().ToString(),
|
||||
EarliestUtc = DateTime.MaxValue,
|
||||
LatestUtc = DateTime.MinValue,
|
||||
NumberOfEvents = 0,
|
||||
NumberOfMoments = 0,
|
||||
};
|
||||
this.log = log;
|
||||
this.workingDir = workingDir;
|
||||
queue.Start();
|
||||
}
|
||||
|
||||
public int NumberOfMoments => reference.NumberOfMoments;
|
||||
|
||||
public void Add(OverwatchMoment moment)
|
||||
{
|
||||
if (moment.Utc < reference.EarliestUtc) reference.EarliestUtc = moment.Utc;
|
||||
if (moment.Utc > reference.LatestUtc) reference.LatestUtc = moment.Utc;
|
||||
reference.NumberOfMoments++;
|
||||
reference.NumberOfEvents += moment.Events.Length;
|
||||
|
||||
var filePath = Path.Combine(workingDir, reference.MomentsFile);
|
||||
|
||||
queue.Add(() =>
|
||||
{
|
||||
File.AppendAllLines(filePath, new[]
|
||||
{
|
||||
Json.Serialize(moment)
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
public OverwatchMomentReference Build()
|
||||
{
|
||||
queue.StopAndJoin();
|
||||
|
||||
log.Debug($"Created reference with {reference.NumberOfMoments} moments and {reference.NumberOfEvents} events...");
|
||||
var result = reference;
|
||||
reference = null!;
|
||||
return result;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,17 @@
|
||||
<Project Sdk="Microsoft.NET.Sdk">
|
||||
|
||||
<PropertyGroup>
|
||||
<TargetFramework>net7.0</TargetFramework>
|
||||
<ImplicitUsings>enable</ImplicitUsings>
|
||||
<Nullable>enable</Nullable>
|
||||
</PropertyGroup>
|
||||
|
||||
<ItemGroup>
|
||||
<PackageReference Include="Newtonsoft.Json" Version="13.0.3" />
|
||||
</ItemGroup>
|
||||
|
||||
<ItemGroup>
|
||||
<ProjectReference Include="..\Logging\Logging.csproj" />
|
||||
</ItemGroup>
|
||||
|
||||
</Project>
|
||||
@@ -0,0 +1,23 @@
|
||||
using Logging;
|
||||
|
||||
namespace OverwatchTranscript
|
||||
{
|
||||
public static class Transcript
|
||||
{
|
||||
public static ITranscriptWriter NewWriter(ILog log)
|
||||
{
|
||||
log = new LogPrefixer(log, "(TranscriptWriter) ");
|
||||
return new TranscriptWriter(log, NewWorkDir());
|
||||
}
|
||||
|
||||
public static ITranscriptReader NewReader(string transcriptFile)
|
||||
{
|
||||
return new TranscriptReader(NewWorkDir(), transcriptFile);
|
||||
}
|
||||
|
||||
private static string NewWorkDir()
|
||||
{
|
||||
return Path.Combine(Path.GetTempPath(), Guid.NewGuid().ToString());
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,8 @@
|
||||
namespace OverwatchTranscript
|
||||
{
|
||||
public static class TranscriptConstants
|
||||
{
|
||||
public const string TranscriptFilename = "transcript.json";
|
||||
public const string ArtifactFolderName = "artifacts";
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,264 @@
|
||||
using Newtonsoft.Json;
|
||||
using System.IO;
|
||||
using System;
|
||||
using System.IO.Compression;
|
||||
using System.Linq;
|
||||
using System.Collections.Generic;
|
||||
using System.Collections.Concurrent;
|
||||
|
||||
namespace OverwatchTranscript
|
||||
{
|
||||
public interface ITranscriptReader
|
||||
{
|
||||
OverwatchCommonHeader Header { get; }
|
||||
T GetHeader<T>(string key);
|
||||
void AddMomentHandler(Action<ActivateMoment> handler);
|
||||
void AddEventHandler<T>(Action<ActivateEvent<T>> handler);
|
||||
bool Next();
|
||||
void Close();
|
||||
}
|
||||
|
||||
public class TranscriptReader : ITranscriptReader
|
||||
{
|
||||
private readonly object handlersLock = new object();
|
||||
private readonly string transcriptFile;
|
||||
private readonly string artifactsFolder;
|
||||
private readonly List<Action<ActivateMoment>> momentHandlers = new List<Action<ActivateMoment>>();
|
||||
private readonly Dictionary<string, List<Action<ActivateMoment, string>>> eventHandlers = new Dictionary<string, List<Action<ActivateMoment, string>>>();
|
||||
private readonly string workingDir;
|
||||
private readonly OverwatchTranscript model;
|
||||
private bool closed;
|
||||
private long momentCounter;
|
||||
private readonly ConcurrentQueue<OverwatchMoment> queue = new ConcurrentQueue<OverwatchMoment>();
|
||||
private readonly Task queueFiller;
|
||||
|
||||
public TranscriptReader(string workingDir, string inputFilename)
|
||||
{
|
||||
closed = false;
|
||||
this.workingDir = workingDir;
|
||||
transcriptFile = Path.Combine(workingDir, TranscriptConstants.TranscriptFilename);
|
||||
artifactsFolder = Path.Combine(workingDir, TranscriptConstants.ArtifactFolderName);
|
||||
|
||||
if (!Directory.Exists(workingDir)) Directory.CreateDirectory(workingDir);
|
||||
if (File.Exists(transcriptFile) || Directory.Exists(artifactsFolder)) throw new Exception("workingdir not clean");
|
||||
|
||||
model = LoadModel(inputFilename);
|
||||
|
||||
queueFiller = Task.Run(() => FillQueue(model, workingDir));
|
||||
}
|
||||
|
||||
public OverwatchCommonHeader Header
|
||||
{
|
||||
get
|
||||
{
|
||||
CheckClosed();
|
||||
return model.Header.Common;
|
||||
}
|
||||
}
|
||||
|
||||
public T GetHeader<T>(string key)
|
||||
{
|
||||
CheckClosed();
|
||||
var value = model.Header.Entries.First(e => e.Key == key).Value;
|
||||
return JsonConvert.DeserializeObject<T>(value)!;
|
||||
}
|
||||
|
||||
public void AddMomentHandler(Action<ActivateMoment> handler)
|
||||
{
|
||||
CheckClosed();
|
||||
lock (handlersLock)
|
||||
{
|
||||
momentHandlers.Add(handler);
|
||||
}
|
||||
}
|
||||
|
||||
public void AddEventHandler<T>(Action<ActivateEvent<T>> handler)
|
||||
{
|
||||
CheckClosed();
|
||||
|
||||
var typeName = typeof(T).FullName;
|
||||
if (string.IsNullOrEmpty(typeName)) throw new Exception("Empty typename for payload");
|
||||
|
||||
lock (handlersLock)
|
||||
{
|
||||
if (eventHandlers.ContainsKey(typeName))
|
||||
{
|
||||
eventHandlers[typeName].Add(CreateEventAction(handler));
|
||||
}
|
||||
else
|
||||
{
|
||||
eventHandlers.Add(typeName, new List<Action<ActivateMoment, string>>
|
||||
{
|
||||
CreateEventAction(handler)
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private readonly object nextLock = new object();
|
||||
private OverwatchMoment? moment = null;
|
||||
private OverwatchMoment? next = null;
|
||||
|
||||
public bool Next()
|
||||
{
|
||||
CheckClosed();
|
||||
|
||||
OverwatchMoment? m = null;
|
||||
TimeSpan? duration = null;
|
||||
lock (nextLock)
|
||||
{
|
||||
if (next == null)
|
||||
{
|
||||
if (!queue.TryDequeue(out moment)) return false;
|
||||
queue.TryDequeue(out next);
|
||||
}
|
||||
else
|
||||
{
|
||||
moment = next;
|
||||
next = null;
|
||||
queue.TryDequeue(out next);
|
||||
}
|
||||
|
||||
m = moment;
|
||||
duration = GetMomentDuration();
|
||||
}
|
||||
|
||||
ActivateMoment(moment, duration);
|
||||
|
||||
return true;
|
||||
}
|
||||
|
||||
public void Close()
|
||||
{
|
||||
CheckClosed();
|
||||
closed = true;
|
||||
|
||||
queueFiller.Wait();
|
||||
|
||||
Directory.Delete(workingDir, true);
|
||||
}
|
||||
|
||||
private Action<ActivateMoment, string> CreateEventAction<T>(Action<ActivateEvent<T>> handler)
|
||||
{
|
||||
return (m, s) =>
|
||||
{
|
||||
handler(new ActivateEvent<T>(m, JsonConvert.DeserializeObject<T>(s)!));
|
||||
};
|
||||
}
|
||||
|
||||
private void FillQueue(OverwatchTranscript model, string workingDir)
|
||||
{
|
||||
var reader = new MomentReader(model, workingDir);
|
||||
|
||||
while (true)
|
||||
{
|
||||
if (closed)
|
||||
{
|
||||
reader.Close();
|
||||
return;
|
||||
}
|
||||
|
||||
while (queue.Count < 10)
|
||||
{
|
||||
var moment = reader.Next();
|
||||
if (moment == null)
|
||||
{
|
||||
reader.Close();
|
||||
return;
|
||||
}
|
||||
queue.Enqueue(moment);
|
||||
}
|
||||
|
||||
Thread.Sleep(1);
|
||||
}
|
||||
}
|
||||
|
||||
private TimeSpan? GetMomentDuration()
|
||||
{
|
||||
if (moment == null) return null;
|
||||
if (next == null) return null;
|
||||
|
||||
return next.Utc - moment.Utc;
|
||||
}
|
||||
|
||||
private void ActivateMoment(OverwatchMoment moment, TimeSpan? duration)
|
||||
{
|
||||
var m = new ActivateMoment(moment.Utc, duration, momentCounter);
|
||||
|
||||
lock (handlersLock)
|
||||
{
|
||||
ActivateMomentHandlers(m);
|
||||
|
||||
foreach (var @event in moment.Events)
|
||||
{
|
||||
ActivateEventHandlers(m, @event);
|
||||
}
|
||||
}
|
||||
|
||||
momentCounter++;
|
||||
}
|
||||
|
||||
private void ActivateMomentHandlers(ActivateMoment m)
|
||||
{
|
||||
foreach (var handler in momentHandlers)
|
||||
{
|
||||
handler(m);
|
||||
}
|
||||
}
|
||||
|
||||
private void ActivateEventHandlers(ActivateMoment m, OverwatchEvent @event)
|
||||
{
|
||||
if (!eventHandlers.ContainsKey(@event.Type)) return;
|
||||
var handlers = eventHandlers[@event.Type];
|
||||
|
||||
foreach (var handler in handlers)
|
||||
{
|
||||
handler(m, @event.Payload);
|
||||
}
|
||||
}
|
||||
|
||||
private OverwatchTranscript LoadModel(string inputFilename)
|
||||
{
|
||||
ZipFile.ExtractToDirectory(inputFilename, workingDir);
|
||||
|
||||
if (!File.Exists(transcriptFile))
|
||||
{
|
||||
closed = true;
|
||||
throw new Exception("Is not a transcript file. Unzipped to: " + workingDir);
|
||||
}
|
||||
|
||||
return JsonConvert.DeserializeObject<OverwatchTranscript>(File.ReadAllText(transcriptFile))!;
|
||||
}
|
||||
|
||||
private void CheckClosed()
|
||||
{
|
||||
if (closed) throw new Exception("Transcript has already been closed.");
|
||||
}
|
||||
}
|
||||
|
||||
public class ActivateMoment
|
||||
{
|
||||
public ActivateMoment(DateTime utc, TimeSpan? duration, long index)
|
||||
{
|
||||
Utc = utc;
|
||||
Duration = duration;
|
||||
Index = index;
|
||||
}
|
||||
|
||||
public DateTime Utc { get; }
|
||||
public TimeSpan? Duration { get; }
|
||||
public long Index { get; }
|
||||
}
|
||||
|
||||
public class ActivateEvent<T>
|
||||
{
|
||||
public ActivateEvent(ActivateMoment moment, T payload)
|
||||
{
|
||||
Moment = moment;
|
||||
Payload = payload;
|
||||
}
|
||||
|
||||
public ActivateMoment Moment { get; }
|
||||
public T Payload { get; }
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,130 @@
|
||||
using Logging;
|
||||
using Newtonsoft.Json;
|
||||
using System.IO.Compression;
|
||||
|
||||
namespace OverwatchTranscript
|
||||
{
|
||||
public interface ITranscriptWriter
|
||||
{
|
||||
void AddHeader(string key, object value);
|
||||
void Add(DateTime utc, object payload);
|
||||
void IncludeArtifact(string filePath);
|
||||
void Write(string outputFilename);
|
||||
}
|
||||
|
||||
public class TranscriptWriter : ITranscriptWriter
|
||||
{
|
||||
private readonly object _lock = new object();
|
||||
private readonly MomentReferenceBuilder builder;
|
||||
private readonly string transcriptFile;
|
||||
private readonly string artifactsFolder;
|
||||
private readonly Dictionary<string, string> header = new Dictionary<string, string>();
|
||||
private readonly BucketSet bucketSet;
|
||||
private readonly ILog log;
|
||||
private readonly string workingDir;
|
||||
private bool closed;
|
||||
|
||||
public TranscriptWriter(ILog log, string workingDir)
|
||||
{
|
||||
closed = false;
|
||||
this.log = log;
|
||||
this.workingDir = workingDir;
|
||||
bucketSet = new BucketSet(log, workingDir);
|
||||
builder = new MomentReferenceBuilder(log, workingDir);
|
||||
transcriptFile = Path.Combine(workingDir, TranscriptConstants.TranscriptFilename);
|
||||
artifactsFolder = Path.Combine(workingDir, TranscriptConstants.ArtifactFolderName);
|
||||
|
||||
if (!Directory.Exists(workingDir)) Directory.CreateDirectory(workingDir);
|
||||
if (File.Exists(transcriptFile) || Directory.Exists(artifactsFolder)) throw new Exception("workingdir not clean");
|
||||
}
|
||||
|
||||
public void Add(DateTime utc, object payload)
|
||||
{
|
||||
CheckClosed();
|
||||
bucketSet.Add(utc, payload);
|
||||
}
|
||||
|
||||
public void AddHeader(string key, object value)
|
||||
{
|
||||
CheckClosed();
|
||||
lock (_lock)
|
||||
{
|
||||
header.Add(key, Json.Serialize(value));
|
||||
}
|
||||
}
|
||||
|
||||
public void IncludeArtifact(string filePath)
|
||||
{
|
||||
CheckClosed();
|
||||
if (!File.Exists(filePath)) throw new Exception("File not found: " + filePath);
|
||||
if (!Directory.Exists(artifactsFolder)) Directory.CreateDirectory(artifactsFolder);
|
||||
var name = Path.GetFileName(filePath);
|
||||
File.Copy(filePath, Path.Combine(artifactsFolder, name), overwrite: false);
|
||||
}
|
||||
|
||||
public void Write(string outputFilename)
|
||||
{
|
||||
CheckClosed();
|
||||
closed = true;
|
||||
|
||||
var momentReferences = builder.Build(bucketSet.FinalizeBuckets());
|
||||
var model = CreateModel(momentReferences);
|
||||
|
||||
File.WriteAllText(transcriptFile, Json.Serialize(model, Formatting.Indented));
|
||||
|
||||
ZipFile.CreateFromDirectory(workingDir, outputFilename);
|
||||
log.Debug($"Transcript written to {outputFilename}");
|
||||
log.Debug($"Common header: {Json.Serialize(model.Header.Common, Formatting.Indented)}");
|
||||
|
||||
Directory.Delete(workingDir, true);
|
||||
log.Debug($"Workdir {workingDir} deleted");
|
||||
}
|
||||
|
||||
private OverwatchTranscript CreateModel(OverwatchMomentReference[] momentReferences)
|
||||
{
|
||||
lock (_lock)
|
||||
{
|
||||
var model = new OverwatchTranscript
|
||||
{
|
||||
Header = new OverwatchHeader
|
||||
{
|
||||
Common = CreateCommonHeader(momentReferences),
|
||||
Entries = header.Select(h =>
|
||||
{
|
||||
return new OverwatchHeaderEntry
|
||||
{
|
||||
Key = h.Key,
|
||||
Value = h.Value
|
||||
};
|
||||
}).ToArray()
|
||||
},
|
||||
MomentReferences = momentReferences
|
||||
};
|
||||
|
||||
header.Clear();
|
||||
return model;
|
||||
}
|
||||
}
|
||||
|
||||
private OverwatchCommonHeader CreateCommonHeader(OverwatchMomentReference[] momentReferences)
|
||||
{
|
||||
var moments = momentReferences.Sum(m => m.NumberOfMoments);
|
||||
var events = momentReferences.Sum(m => m.NumberOfEvents);
|
||||
var earliest = momentReferences.Min(m => m.EarliestUtc);
|
||||
var latest = momentReferences.Max(m => m.LatestUtc);
|
||||
|
||||
return new OverwatchCommonHeader
|
||||
{
|
||||
NumberOfMoments = moments,
|
||||
NumberOfEvents = events,
|
||||
EarliestUtc = earliest,
|
||||
LatestUtc = latest
|
||||
};
|
||||
}
|
||||
|
||||
private void CheckClosed()
|
||||
{
|
||||
if (closed) throw new Exception("Transcript has already been written. Cannot modify or write again.");
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -2,7 +2,7 @@
|
||||
{
|
||||
public class BlockInterval
|
||||
{
|
||||
public BlockInterval(ulong from, ulong to)
|
||||
public BlockInterval(TimeRange timeRange, ulong from, ulong to)
|
||||
{
|
||||
if (from < to)
|
||||
{
|
||||
@@ -14,10 +14,13 @@
|
||||
From = to;
|
||||
To = from;
|
||||
}
|
||||
TimeRange = timeRange;
|
||||
}
|
||||
|
||||
public ulong From { get; }
|
||||
public ulong To { get; }
|
||||
public TimeRange TimeRange { get; }
|
||||
public ulong NumberOfBlocks => To - From;
|
||||
|
||||
public override string ToString()
|
||||
{
|
||||
|
||||
@@ -13,7 +13,6 @@
|
||||
|
||||
public long SizeInBytes { get; }
|
||||
|
||||
|
||||
public long ToMB()
|
||||
{
|
||||
return SizeInBytes / (1024 * 1024);
|
||||
|
||||
@@ -11,5 +11,16 @@
|
||||
remainingItems.RemoveAt(i);
|
||||
return result;
|
||||
}
|
||||
|
||||
public static T[] Shuffled<T>(T[] items)
|
||||
{
|
||||
var result = new List<T>();
|
||||
var source = items.ToList();
|
||||
while (source.Any())
|
||||
{
|
||||
result.Add(RandomUtils.PickOneRandom(source));
|
||||
}
|
||||
return result.ToArray();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,28 @@
|
||||
namespace Utils
|
||||
{
|
||||
public static class RollingAverage
|
||||
{
|
||||
/// <param name="currentAverage">Value of average before new value is added.</param>
|
||||
/// <param name="newNumberOfValues">Number of values in average after new value is added.</param>
|
||||
/// <param name="newValue">New value to be added.</param>
|
||||
/// <returns>New average value.</returns>
|
||||
/// <exception cref="Exception">newNumberOfValues must be 1 or greater.</exception>
|
||||
public static float GetNewAverage(float currentAverage, int newNumberOfValues, float newValue)
|
||||
{
|
||||
if (newNumberOfValues < 1) throw new Exception("Should be at least 1 value.");
|
||||
|
||||
float n = newNumberOfValues;
|
||||
var originalValue = currentAverage;
|
||||
var originalValueWeight = ((n - 1.0f) / n);
|
||||
var newValueWeight = (1.0f / n);
|
||||
return GetWeightedAverage(originalValue, originalValueWeight, newValue, newValueWeight);
|
||||
}
|
||||
|
||||
public static float GetWeightedAverage(float value1, float weight1, float value2, float weight2)
|
||||
{
|
||||
float totalWeight = weight1 + weight2;
|
||||
if (totalWeight == 0.0f) return 0.0f;
|
||||
return ((value1 * weight1) + (value2 * weight2)) / totalWeight;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,13 @@
|
||||
namespace Utils
|
||||
{
|
||||
public static class Str
|
||||
{
|
||||
public static string Between(string input, string open, string close)
|
||||
{
|
||||
var openIndex = input.IndexOf(open) + open.Length;
|
||||
var closeIndex = input.LastIndexOf(close);
|
||||
|
||||
return input.Substring(openIndex, closeIndex - openIndex);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,73 @@
|
||||
using CodexContractsPlugin.Marketplace;
|
||||
using Utils;
|
||||
|
||||
namespace CodexContractsPlugin.ChainMonitor
|
||||
{
|
||||
public class ChainEvents
|
||||
{
|
||||
private ChainEvents(
|
||||
BlockInterval blockInterval,
|
||||
Request[] requests,
|
||||
RequestFulfilledEventDTO[] fulfilled,
|
||||
RequestCancelledEventDTO[] cancelled,
|
||||
RequestFailedEventDTO[] failed,
|
||||
SlotFilledEventDTO[] slotFilled,
|
||||
SlotFreedEventDTO[] slotFreed
|
||||
)
|
||||
{
|
||||
BlockInterval = blockInterval;
|
||||
Requests = requests;
|
||||
Fulfilled = fulfilled;
|
||||
Cancelled = cancelled;
|
||||
Failed = failed;
|
||||
SlotFilled = slotFilled;
|
||||
SlotFreed = slotFreed;
|
||||
}
|
||||
|
||||
public BlockInterval BlockInterval { get; }
|
||||
public Request[] Requests { get; }
|
||||
public RequestFulfilledEventDTO[] Fulfilled { get; }
|
||||
public RequestCancelledEventDTO[] Cancelled { get; }
|
||||
public RequestFailedEventDTO[] Failed { get; }
|
||||
public SlotFilledEventDTO[] SlotFilled { get; }
|
||||
public SlotFreedEventDTO[] SlotFreed { get; }
|
||||
|
||||
public IHasBlock[] All
|
||||
{
|
||||
get
|
||||
{
|
||||
var all = new List<IHasBlock>();
|
||||
all.AddRange(Requests);
|
||||
all.AddRange(Fulfilled);
|
||||
all.AddRange(Cancelled);
|
||||
all.AddRange(Failed);
|
||||
all.AddRange(SlotFilled);
|
||||
all.AddRange(SlotFreed);
|
||||
return all.ToArray();
|
||||
}
|
||||
}
|
||||
|
||||
public static ChainEvents FromBlockInterval(ICodexContracts contracts, BlockInterval blockInterval)
|
||||
{
|
||||
return FromContractEvents(contracts.GetEvents(blockInterval));
|
||||
}
|
||||
|
||||
public static ChainEvents FromTimeRange(ICodexContracts contracts, TimeRange timeRange)
|
||||
{
|
||||
return FromContractEvents(contracts.GetEvents(timeRange));
|
||||
}
|
||||
|
||||
public static ChainEvents FromContractEvents(ICodexContractsEvents events)
|
||||
{
|
||||
return new ChainEvents(
|
||||
events.BlockInterval,
|
||||
events.GetStorageRequests(),
|
||||
events.GetRequestFulfilledEvents(),
|
||||
events.GetRequestCancelledEvents(),
|
||||
events.GetRequestFailedEvents(),
|
||||
events.GetSlotFilledEvents(),
|
||||
events.GetSlotFreedEvents()
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,176 @@
|
||||
using CodexContractsPlugin.Marketplace;
|
||||
using GethPlugin;
|
||||
using Logging;
|
||||
using NethereumWorkflow.BlockUtils;
|
||||
using System.Numerics;
|
||||
using Utils;
|
||||
|
||||
namespace CodexContractsPlugin.ChainMonitor
|
||||
{
|
||||
public interface IChainStateChangeHandler
|
||||
{
|
||||
void OnNewRequest(RequestEvent requestEvent);
|
||||
void OnRequestFinished(RequestEvent requestEvent);
|
||||
void OnRequestFulfilled(RequestEvent requestEvent);
|
||||
void OnRequestCancelled(RequestEvent requestEvent);
|
||||
void OnRequestFailed(RequestEvent requestEvent);
|
||||
void OnSlotFilled(RequestEvent requestEvent, EthAddress host, BigInteger slotIndex);
|
||||
void OnSlotFreed(RequestEvent requestEvent, BigInteger slotIndex);
|
||||
}
|
||||
|
||||
public class RequestEvent
|
||||
{
|
||||
public RequestEvent(BlockTimeEntry block, IChainStateRequest request)
|
||||
{
|
||||
Block = block;
|
||||
Request = request;
|
||||
}
|
||||
|
||||
public BlockTimeEntry Block { get; }
|
||||
public IChainStateRequest Request { get; }
|
||||
}
|
||||
|
||||
public class ChainState
|
||||
{
|
||||
private readonly List<ChainStateRequest> requests = new List<ChainStateRequest>();
|
||||
private readonly ILog log;
|
||||
private readonly ICodexContracts contracts;
|
||||
private readonly IChainStateChangeHandler handler;
|
||||
|
||||
public ChainState(ILog log, ICodexContracts contracts, IChainStateChangeHandler changeHandler, DateTime startUtc)
|
||||
{
|
||||
this.log = new LogPrefixer(log, "(ChainState) ");
|
||||
this.contracts = contracts;
|
||||
handler = changeHandler;
|
||||
TotalSpan = new TimeRange(startUtc, startUtc);
|
||||
}
|
||||
|
||||
public TimeRange TotalSpan { get; private set; }
|
||||
public IChainStateRequest[] Requests => requests.ToArray();
|
||||
|
||||
public void Update()
|
||||
{
|
||||
Update(DateTime.UtcNow);
|
||||
}
|
||||
|
||||
public void Update(DateTime toUtc)
|
||||
{
|
||||
var span = new TimeRange(TotalSpan.To, toUtc);
|
||||
var events = ChainEvents.FromTimeRange(contracts, span);
|
||||
Apply(events);
|
||||
|
||||
TotalSpan = new TimeRange(TotalSpan.From, span.To);
|
||||
}
|
||||
|
||||
private void Apply(ChainEvents events)
|
||||
{
|
||||
if (events.BlockInterval.TimeRange.From < TotalSpan.From)
|
||||
throw new Exception("Attempt to update ChainState with set of events from before its current record.");
|
||||
|
||||
log.Log($"ChainState updating: {events.BlockInterval}");
|
||||
|
||||
// Run through each block and apply the events to the state in order.
|
||||
var span = events.BlockInterval.TimeRange.Duration;
|
||||
var numBlocks = events.BlockInterval.NumberOfBlocks;
|
||||
var spanPerBlock = span / numBlocks;
|
||||
|
||||
var eventUtc = events.BlockInterval.TimeRange.From;
|
||||
for (var b = events.BlockInterval.From; b <= events.BlockInterval.To; b++)
|
||||
{
|
||||
var blockEvents = events.All.Where(e => e.Block.BlockNumber == b).ToArray();
|
||||
ApplyEvents(b, blockEvents, eventUtc);
|
||||
|
||||
eventUtc += spanPerBlock;
|
||||
}
|
||||
}
|
||||
|
||||
private void ApplyEvents(ulong blockNumber, IHasBlock[] blockEvents, DateTime eventsUtc)
|
||||
{
|
||||
foreach (var e in blockEvents)
|
||||
{
|
||||
dynamic d = e;
|
||||
ApplyEvent(d);
|
||||
}
|
||||
|
||||
ApplyTimeImplicitEvents(blockNumber, eventsUtc);
|
||||
}
|
||||
|
||||
private void ApplyEvent(Request request)
|
||||
{
|
||||
if (requests.Any(r => Equal(r.Request.RequestId, request.RequestId)))
|
||||
throw new Exception("Received NewRequest event for id that already exists.");
|
||||
|
||||
var newRequest = new ChainStateRequest(log, request, RequestState.New);
|
||||
requests.Add(newRequest);
|
||||
|
||||
handler.OnNewRequest(new RequestEvent(request.Block, newRequest));
|
||||
}
|
||||
|
||||
private void ApplyEvent(RequestFulfilledEventDTO @event)
|
||||
{
|
||||
var r = FindRequest(@event.RequestId);
|
||||
if (r == null) return;
|
||||
r.UpdateState(@event.Block.BlockNumber, RequestState.Started);
|
||||
handler.OnRequestFulfilled(new RequestEvent(@event.Block, r));
|
||||
}
|
||||
|
||||
private void ApplyEvent(RequestCancelledEventDTO @event)
|
||||
{
|
||||
var r = FindRequest(@event.RequestId);
|
||||
if (r == null) return;
|
||||
r.UpdateState(@event.Block.BlockNumber, RequestState.Cancelled);
|
||||
handler.OnRequestCancelled(new RequestEvent(@event.Block, r));
|
||||
}
|
||||
|
||||
private void ApplyEvent(RequestFailedEventDTO @event)
|
||||
{
|
||||
var r = FindRequest(@event.RequestId);
|
||||
if (r == null) return;
|
||||
r.UpdateState(@event.Block.BlockNumber, RequestState.Failed);
|
||||
handler.OnRequestFailed(new RequestEvent(@event.Block, r));
|
||||
}
|
||||
|
||||
private void ApplyEvent(SlotFilledEventDTO @event)
|
||||
{
|
||||
var r = FindRequest(@event.RequestId);
|
||||
if (r == null) return;
|
||||
r.Hosts.Add(@event.Host, (int)@event.SlotIndex);
|
||||
r.Log($"[{@event.Block.BlockNumber}] SlotFilled (host:'{@event.Host}', slotIndex:{@event.SlotIndex})");
|
||||
handler.OnSlotFilled(new RequestEvent(@event.Block, r), @event.Host, @event.SlotIndex);
|
||||
}
|
||||
|
||||
private void ApplyEvent(SlotFreedEventDTO @event)
|
||||
{
|
||||
var r = FindRequest(@event.RequestId);
|
||||
if (r == null) return;
|
||||
r.Hosts.RemoveHost((int)@event.SlotIndex);
|
||||
r.Log($"[{@event.Block.BlockNumber}] SlotFreed (slotIndex:{@event.SlotIndex})");
|
||||
handler.OnSlotFreed(new RequestEvent(@event.Block, r), @event.SlotIndex);
|
||||
}
|
||||
|
||||
private void ApplyTimeImplicitEvents(ulong blockNumber, DateTime eventsUtc)
|
||||
{
|
||||
foreach (var r in requests)
|
||||
{
|
||||
if (r.State == RequestState.Started
|
||||
&& r.FinishedUtc < eventsUtc)
|
||||
{
|
||||
r.UpdateState(blockNumber, RequestState.Finished);
|
||||
handler.OnRequestFinished(new RequestEvent(new BlockTimeEntry(blockNumber, eventsUtc), r));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private ChainStateRequest? FindRequest(byte[] requestId)
|
||||
{
|
||||
var r = requests.SingleOrDefault(r => Equal(r.Request.RequestId, requestId));
|
||||
if (r == null) log.Log("Unable to find request by ID!");
|
||||
return r;
|
||||
}
|
||||
|
||||
private bool Equal(byte[] a, byte[] b)
|
||||
{
|
||||
return a.SequenceEqual(b);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,55 @@
|
||||
using GethPlugin;
|
||||
using System;
|
||||
using System.Collections.Generic;
|
||||
using System.Linq;
|
||||
using System.Numerics;
|
||||
using System.Text;
|
||||
using System.Threading.Tasks;
|
||||
|
||||
namespace CodexContractsPlugin.ChainMonitor
|
||||
{
|
||||
public class ChainStateChangeHandlerMux : IChainStateChangeHandler
|
||||
{
|
||||
public ChainStateChangeHandlerMux(params IChainStateChangeHandler[] handlers)
|
||||
{
|
||||
Handlers = handlers.ToList();
|
||||
}
|
||||
|
||||
public List<IChainStateChangeHandler> Handlers { get; } = new List<IChainStateChangeHandler>();
|
||||
|
||||
public void OnNewRequest(RequestEvent requestEvent)
|
||||
{
|
||||
foreach (var handler in Handlers) handler.OnNewRequest(requestEvent);
|
||||
}
|
||||
|
||||
public void OnRequestCancelled(RequestEvent requestEvent)
|
||||
{
|
||||
foreach (var handler in Handlers) handler.OnRequestCancelled(requestEvent);
|
||||
}
|
||||
|
||||
public void OnRequestFailed(RequestEvent requestEvent)
|
||||
{
|
||||
foreach (var handler in Handlers) handler.OnRequestFailed(requestEvent);
|
||||
}
|
||||
|
||||
public void OnRequestFinished(RequestEvent requestEvent)
|
||||
{
|
||||
foreach (var handler in Handlers) handler.OnRequestFinished(requestEvent);
|
||||
}
|
||||
|
||||
public void OnRequestFulfilled(RequestEvent requestEvent)
|
||||
{
|
||||
foreach (var handler in Handlers) handler.OnRequestFulfilled(requestEvent);
|
||||
}
|
||||
|
||||
public void OnSlotFilled(RequestEvent requestEvent, EthAddress host, BigInteger slotIndex)
|
||||
{
|
||||
foreach (var handler in Handlers) handler.OnSlotFilled(requestEvent, host, slotIndex);
|
||||
}
|
||||
|
||||
public void OnSlotFreed(RequestEvent requestEvent, BigInteger slotIndex)
|
||||
{
|
||||
foreach (var handler in Handlers) handler.OnSlotFreed(requestEvent, slotIndex);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,80 @@
|
||||
using CodexContractsPlugin.Marketplace;
|
||||
using GethPlugin;
|
||||
using Logging;
|
||||
|
||||
namespace CodexContractsPlugin.ChainMonitor
|
||||
{
|
||||
public interface IChainStateRequest
|
||||
{
|
||||
Request Request { get; }
|
||||
RequestState State { get; }
|
||||
DateTime ExpiryUtc { get; }
|
||||
DateTime FinishedUtc { get; }
|
||||
EthAddress Client { get; }
|
||||
RequestHosts Hosts { get; }
|
||||
}
|
||||
|
||||
public class ChainStateRequest : IChainStateRequest
|
||||
{
|
||||
private readonly ILog log;
|
||||
|
||||
public ChainStateRequest(ILog log, Request request, RequestState state)
|
||||
{
|
||||
this.log = log;
|
||||
Request = request;
|
||||
State = state;
|
||||
|
||||
ExpiryUtc = request.Block.Utc + TimeSpan.FromSeconds((double)request.Expiry);
|
||||
FinishedUtc = request.Block.Utc + TimeSpan.FromSeconds((double)request.Ask.Duration);
|
||||
|
||||
Log($"[{request.Block.BlockNumber}] Created as {State}.");
|
||||
|
||||
Client = new EthAddress(request.Client);
|
||||
Hosts = new RequestHosts();
|
||||
}
|
||||
|
||||
public Request Request { get; }
|
||||
public RequestState State { get; private set; }
|
||||
public DateTime ExpiryUtc { get; }
|
||||
public DateTime FinishedUtc { get; }
|
||||
public EthAddress Client { get; }
|
||||
public RequestHosts Hosts { get; }
|
||||
|
||||
public void UpdateState(ulong blockNumber, RequestState newState)
|
||||
{
|
||||
Log($"[{blockNumber}] Transit: {State} -> {newState}");
|
||||
State = newState;
|
||||
}
|
||||
|
||||
public void Log(string msg)
|
||||
{
|
||||
log.Log($"Request '{Request.Id}': {msg}");
|
||||
}
|
||||
}
|
||||
|
||||
public class RequestHosts
|
||||
{
|
||||
private readonly Dictionary<int, EthAddress> hosts = new Dictionary<int, EthAddress>();
|
||||
|
||||
public void Add(EthAddress host, int index)
|
||||
{
|
||||
hosts.Add(index, host);
|
||||
}
|
||||
|
||||
public void RemoveHost(int index)
|
||||
{
|
||||
hosts.Remove(index);
|
||||
}
|
||||
|
||||
public EthAddress? GetHost(int index)
|
||||
{
|
||||
if (!hosts.ContainsKey(index)) return null;
|
||||
return hosts[index];
|
||||
}
|
||||
|
||||
public EthAddress[] GetHosts()
|
||||
{
|
||||
return hosts.Values.ToArray();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,36 @@
|
||||
using GethPlugin;
|
||||
using System.Numerics;
|
||||
|
||||
namespace CodexContractsPlugin.ChainMonitor
|
||||
{
|
||||
public class DoNothingChainEventHandler : IChainStateChangeHandler
|
||||
{
|
||||
public void OnNewRequest(RequestEvent requestEvent)
|
||||
{
|
||||
}
|
||||
|
||||
public void OnRequestCancelled(RequestEvent requestEvent)
|
||||
{
|
||||
}
|
||||
|
||||
public void OnRequestFailed(RequestEvent requestEvent)
|
||||
{
|
||||
}
|
||||
|
||||
public void OnRequestFinished(RequestEvent requestEvent)
|
||||
{
|
||||
}
|
||||
|
||||
public void OnRequestFulfilled(RequestEvent requestEvent)
|
||||
{
|
||||
}
|
||||
|
||||
public void OnSlotFilled(RequestEvent requestEvent, EthAddress host, BigInteger slotIndex)
|
||||
{
|
||||
}
|
||||
|
||||
public void OnSlotFreed(RequestEvent requestEvent, BigInteger slotIndex)
|
||||
{
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -2,10 +2,8 @@
|
||||
using GethPlugin;
|
||||
using Logging;
|
||||
using Nethereum.ABI;
|
||||
using Nethereum.Hex.HexTypes;
|
||||
using Nethereum.Util;
|
||||
using NethereumWorkflow;
|
||||
using NethereumWorkflow.BlockUtils;
|
||||
using Newtonsoft.Json;
|
||||
using Newtonsoft.Json.Converters;
|
||||
using Utils;
|
||||
@@ -22,13 +20,10 @@ namespace CodexContractsPlugin
|
||||
TestToken GetTestTokenBalance(IHasEthAddress owner);
|
||||
TestToken GetTestTokenBalance(EthAddress ethAddress);
|
||||
|
||||
Request[] GetStorageRequests(BlockInterval blockRange);
|
||||
ICodexContractsEvents GetEvents(TimeRange timeRange);
|
||||
ICodexContractsEvents GetEvents(BlockInterval blockInterval);
|
||||
EthAddress? GetSlotHost(Request storageRequest, decimal slotIndex);
|
||||
RequestState GetRequestState(Request request);
|
||||
RequestFulfilledEventDTO[] GetRequestFulfilledEvents(BlockInterval blockRange);
|
||||
RequestCancelledEventDTO[] GetRequestCancelledEvents(BlockInterval blockRange);
|
||||
SlotFilledEventDTO[] GetSlotFilledEvents(BlockInterval blockRange);
|
||||
SlotFreedEventDTO[] GetSlotFreedEvents(BlockInterval blockRange);
|
||||
}
|
||||
|
||||
[JsonConverter(typeof(StringEnumConverter))]
|
||||
@@ -81,65 +76,14 @@ namespace CodexContractsPlugin
|
||||
return balance.TstWei();
|
||||
}
|
||||
|
||||
public Request[] GetStorageRequests(BlockInterval blockRange)
|
||||
public ICodexContractsEvents GetEvents(TimeRange timeRange)
|
||||
{
|
||||
var events = gethNode.GetEvents<StorageRequestedEventDTO>(Deployment.MarketplaceAddress, blockRange);
|
||||
var i = StartInteraction();
|
||||
return events
|
||||
.Select(e =>
|
||||
{
|
||||
var requestEvent = i.GetRequest(Deployment.MarketplaceAddress, e.Event.RequestId);
|
||||
var result = requestEvent.ReturnValue1;
|
||||
result.Block = GetBlock(e.Log.BlockNumber.ToUlong());
|
||||
result.RequestId = e.Event.RequestId;
|
||||
return result;
|
||||
})
|
||||
.ToArray();
|
||||
return GetEvents(gethNode.ConvertTimeRangeToBlockRange(timeRange));
|
||||
}
|
||||
|
||||
public RequestFulfilledEventDTO[] GetRequestFulfilledEvents(BlockInterval blockRange)
|
||||
public ICodexContractsEvents GetEvents(BlockInterval blockInterval)
|
||||
{
|
||||
var events = gethNode.GetEvents<RequestFulfilledEventDTO>(Deployment.MarketplaceAddress, blockRange);
|
||||
return events.Select(e =>
|
||||
{
|
||||
var result = e.Event;
|
||||
result.Block = GetBlock(e.Log.BlockNumber.ToUlong());
|
||||
return result;
|
||||
}).ToArray();
|
||||
}
|
||||
|
||||
public RequestCancelledEventDTO[] GetRequestCancelledEvents(BlockInterval blockRange)
|
||||
{
|
||||
var events = gethNode.GetEvents<RequestCancelledEventDTO>(Deployment.MarketplaceAddress, blockRange);
|
||||
return events.Select(e =>
|
||||
{
|
||||
var result = e.Event;
|
||||
result.Block = GetBlock(e.Log.BlockNumber.ToUlong());
|
||||
return result;
|
||||
}).ToArray();
|
||||
}
|
||||
|
||||
public SlotFilledEventDTO[] GetSlotFilledEvents(BlockInterval blockRange)
|
||||
{
|
||||
var events = gethNode.GetEvents<SlotFilledEventDTO>(Deployment.MarketplaceAddress, blockRange);
|
||||
return events.Select(e =>
|
||||
{
|
||||
var result = e.Event;
|
||||
result.Block = GetBlock(e.Log.BlockNumber.ToUlong());
|
||||
result.Host = GetEthAddressFromTransaction(e.Log.TransactionHash);
|
||||
return result;
|
||||
}).ToArray();
|
||||
}
|
||||
|
||||
public SlotFreedEventDTO[] GetSlotFreedEvents(BlockInterval blockRange)
|
||||
{
|
||||
var events = gethNode.GetEvents<SlotFreedEventDTO>(Deployment.MarketplaceAddress, blockRange);
|
||||
return events.Select(e =>
|
||||
{
|
||||
var result = e.Event;
|
||||
result.Block = GetBlock(e.Log.BlockNumber.ToUlong());
|
||||
return result;
|
||||
}).ToArray();
|
||||
return new CodexContractsEvents(log, gethNode, Deployment, blockInterval);
|
||||
}
|
||||
|
||||
public EthAddress? GetSlotHost(Request storageRequest, decimal slotIndex)
|
||||
@@ -170,17 +114,6 @@ namespace CodexContractsPlugin
|
||||
return gethNode.Call<RequestStateFunction, RequestState>(Deployment.MarketplaceAddress, func);
|
||||
}
|
||||
|
||||
private BlockTimeEntry GetBlock(ulong number)
|
||||
{
|
||||
return gethNode.GetBlockForNumber(number);
|
||||
}
|
||||
|
||||
private EthAddress GetEthAddressFromTransaction(string transactionHash)
|
||||
{
|
||||
var transaction = gethNode.GetTransaction(transactionHash);
|
||||
return new EthAddress(transaction.From);
|
||||
}
|
||||
|
||||
private ContractInteractions StartInteraction()
|
||||
{
|
||||
return new ContractInteractions(log, gethNode);
|
||||
|
||||
@@ -0,0 +1,120 @@
|
||||
using CodexContractsPlugin.Marketplace;
|
||||
using GethPlugin;
|
||||
using Logging;
|
||||
using Nethereum.Hex.HexTypes;
|
||||
using NethereumWorkflow.BlockUtils;
|
||||
using Utils;
|
||||
|
||||
namespace CodexContractsPlugin
|
||||
{
|
||||
public interface ICodexContractsEvents
|
||||
{
|
||||
BlockInterval BlockInterval { get; }
|
||||
Request[] GetStorageRequests();
|
||||
RequestFulfilledEventDTO[] GetRequestFulfilledEvents();
|
||||
RequestCancelledEventDTO[] GetRequestCancelledEvents();
|
||||
RequestFailedEventDTO[] GetRequestFailedEvents();
|
||||
SlotFilledEventDTO[] GetSlotFilledEvents();
|
||||
SlotFreedEventDTO[] GetSlotFreedEvents();
|
||||
}
|
||||
|
||||
public class CodexContractsEvents : ICodexContractsEvents
|
||||
{
|
||||
private readonly ILog log;
|
||||
private readonly IGethNode gethNode;
|
||||
private readonly CodexContractsDeployment deployment;
|
||||
|
||||
public CodexContractsEvents(ILog log, IGethNode gethNode, CodexContractsDeployment deployment, BlockInterval blockInterval)
|
||||
{
|
||||
this.log = log;
|
||||
this.gethNode = gethNode;
|
||||
this.deployment = deployment;
|
||||
BlockInterval = blockInterval;
|
||||
}
|
||||
|
||||
public BlockInterval BlockInterval { get; }
|
||||
|
||||
public Request[] GetStorageRequests()
|
||||
{
|
||||
var events = gethNode.GetEvents<StorageRequestedEventDTO>(deployment.MarketplaceAddress, BlockInterval);
|
||||
var i = new ContractInteractions(log, gethNode);
|
||||
return events
|
||||
.Select(e =>
|
||||
{
|
||||
var requestEvent = i.GetRequest(deployment.MarketplaceAddress, e.Event.RequestId);
|
||||
var result = requestEvent.ReturnValue1;
|
||||
result.Block = GetBlock(e.Log.BlockNumber.ToUlong());
|
||||
result.RequestId = e.Event.RequestId;
|
||||
return result;
|
||||
})
|
||||
.ToArray();
|
||||
}
|
||||
|
||||
public RequestFulfilledEventDTO[] GetRequestFulfilledEvents()
|
||||
{
|
||||
var events = gethNode.GetEvents<RequestFulfilledEventDTO>(deployment.MarketplaceAddress, BlockInterval);
|
||||
return events.Select(e =>
|
||||
{
|
||||
var result = e.Event;
|
||||
result.Block = GetBlock(e.Log.BlockNumber.ToUlong());
|
||||
return result;
|
||||
}).ToArray();
|
||||
}
|
||||
|
||||
public RequestCancelledEventDTO[] GetRequestCancelledEvents()
|
||||
{
|
||||
var events = gethNode.GetEvents<RequestCancelledEventDTO>(deployment.MarketplaceAddress, BlockInterval);
|
||||
return events.Select(e =>
|
||||
{
|
||||
var result = e.Event;
|
||||
result.Block = GetBlock(e.Log.BlockNumber.ToUlong());
|
||||
return result;
|
||||
}).ToArray();
|
||||
}
|
||||
|
||||
public RequestFailedEventDTO[] GetRequestFailedEvents()
|
||||
{
|
||||
var events = gethNode.GetEvents<RequestFailedEventDTO>(deployment.MarketplaceAddress, BlockInterval);
|
||||
return events.Select(e =>
|
||||
{
|
||||
var result = e.Event;
|
||||
result.Block = GetBlock(e.Log.BlockNumber.ToUlong());
|
||||
return result;
|
||||
}).ToArray();
|
||||
}
|
||||
|
||||
public SlotFilledEventDTO[] GetSlotFilledEvents()
|
||||
{
|
||||
var events = gethNode.GetEvents<SlotFilledEventDTO>(deployment.MarketplaceAddress, BlockInterval);
|
||||
return events.Select(e =>
|
||||
{
|
||||
var result = e.Event;
|
||||
result.Block = GetBlock(e.Log.BlockNumber.ToUlong());
|
||||
result.Host = GetEthAddressFromTransaction(e.Log.TransactionHash);
|
||||
return result;
|
||||
}).ToArray();
|
||||
}
|
||||
|
||||
public SlotFreedEventDTO[] GetSlotFreedEvents()
|
||||
{
|
||||
var events = gethNode.GetEvents<SlotFreedEventDTO>(deployment.MarketplaceAddress, BlockInterval);
|
||||
return events.Select(e =>
|
||||
{
|
||||
var result = e.Event;
|
||||
result.Block = GetBlock(e.Log.BlockNumber.ToUlong());
|
||||
return result;
|
||||
}).ToArray();
|
||||
}
|
||||
|
||||
private BlockTimeEntry GetBlock(ulong number)
|
||||
{
|
||||
return gethNode.GetBlockForNumber(number);
|
||||
}
|
||||
|
||||
private EthAddress GetEthAddressFromTransaction(string transactionHash)
|
||||
{
|
||||
var transaction = gethNode.GetTransaction(transactionHash);
|
||||
return new EthAddress(transaction.From);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -6,6 +6,11 @@
|
||||
<Nullable>enable</Nullable>
|
||||
</PropertyGroup>
|
||||
|
||||
<ItemGroup>
|
||||
<PackageReference Include="Nethereum.Generators" Version="4.21.4" />
|
||||
<PackageReference Include="Nethereum.Generators.Net" Version="4.21.4" />
|
||||
</ItemGroup>
|
||||
|
||||
<ItemGroup>
|
||||
<ProjectReference Include="..\..\Framework\Core\Core.csproj" />
|
||||
<ProjectReference Include="..\GethPlugin\GethPlugin.csproj" />
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
using Core;
|
||||
using CodexContractsPlugin.Marketplace;
|
||||
using Core;
|
||||
using GethPlugin;
|
||||
using KubernetesWorkflow;
|
||||
using KubernetesWorkflow.Types;
|
||||
@@ -64,7 +65,8 @@ namespace CodexContractsPlugin
|
||||
|
||||
var extractor = new ContractsContainerInfoExtractor(tools.GetLog(), workflow, container);
|
||||
var marketplaceAddress = extractor.ExtractMarketplaceAddress();
|
||||
var abi = extractor.ExtractMarketplaceAbi();
|
||||
var (abi, bytecode) = extractor.ExtractMarketplaceAbiAndByteCode();
|
||||
EnsureCompatbility(abi, bytecode);
|
||||
|
||||
var interaction = new ContractInteractions(tools.GetLog(), gethNode);
|
||||
var tokenAddress = interaction.GetTokenAddress(marketplaceAddress);
|
||||
@@ -78,6 +80,18 @@ namespace CodexContractsPlugin
|
||||
return new CodexContractsDeployment(marketplaceAddress, abi, tokenAddress);
|
||||
}
|
||||
|
||||
private void EnsureCompatbility(string abi, string bytecode)
|
||||
{
|
||||
var expectedByteCode = MarketplaceDeploymentBase.BYTECODE.ToLowerInvariant();
|
||||
|
||||
if (bytecode != expectedByteCode)
|
||||
{
|
||||
Log("Deployed contract is incompatible with current build of CodexContracts plugin. Running self-updater...");
|
||||
var selfUpdater = new SelfUpdater();
|
||||
selfUpdater.Update(abi, bytecode);
|
||||
}
|
||||
}
|
||||
|
||||
private void Log(string msg)
|
||||
{
|
||||
tools.GetLog().Log(msg);
|
||||
|
||||
@@ -31,14 +31,14 @@ namespace CodexContractsPlugin
|
||||
return marketplaceAddress;
|
||||
}
|
||||
|
||||
public string ExtractMarketplaceAbi()
|
||||
public (string, string) ExtractMarketplaceAbiAndByteCode()
|
||||
{
|
||||
log.Debug();
|
||||
var marketplaceAbi = Retry(FetchMarketplaceAbi);
|
||||
if (string.IsNullOrEmpty(marketplaceAbi)) throw new InvalidOperationException("Unable to fetch marketplace artifacts from codex-contracts node. Test infra failure.");
|
||||
var (abi, bytecode) = Retry(FetchMarketplaceAbiAndByteCode);
|
||||
if (string.IsNullOrEmpty(abi)) throw new InvalidOperationException("Unable to fetch marketplace artifacts from codex-contracts node. Test infra failure.");
|
||||
|
||||
log.Debug("Got Marketplace ABI: " + marketplaceAbi);
|
||||
return marketplaceAbi;
|
||||
log.Debug("Got Marketplace ABI: " + abi);
|
||||
return (abi, bytecode);
|
||||
}
|
||||
|
||||
private string FetchMarketplaceAddress()
|
||||
@@ -48,7 +48,7 @@ namespace CodexContractsPlugin
|
||||
return marketplace!.address;
|
||||
}
|
||||
|
||||
private string FetchMarketplaceAbi()
|
||||
private (string, string) FetchMarketplaceAbiAndByteCode()
|
||||
{
|
||||
var json = workflow.ExecuteCommand(container, "cat", CodexContractsContainerRecipe.MarketplaceArtifactFilename);
|
||||
|
||||
@@ -56,19 +56,12 @@ namespace CodexContractsPlugin
|
||||
var abi = artifact["abi"];
|
||||
var byteCode = artifact["bytecode"];
|
||||
var abiResult = abi!.ToString(Formatting.None);
|
||||
var byteCodeResult = byteCode!.ToString(Formatting.None);
|
||||
|
||||
if (byteCodeResult
|
||||
.ToLowerInvariant()
|
||||
.Replace("\"", "") != MarketplaceDeploymentBase.BYTECODE.ToLowerInvariant())
|
||||
{
|
||||
throw new Exception("BYTECODE in CodexContractsPlugin does not match BYTECODE deployed by container. Update Marketplace.cs generated code?");
|
||||
}
|
||||
|
||||
return abiResult;
|
||||
var byteCodeResult = byteCode!.ToString(Formatting.None).ToLowerInvariant().Replace("\"", "");
|
||||
|
||||
return (abiResult, byteCodeResult);
|
||||
}
|
||||
|
||||
private static string Retry(Func<string> fetch)
|
||||
private static T Retry<T>(Func<T> fetch)
|
||||
{
|
||||
return Time.Retry(fetch, nameof(ContractsContainerInfoExtractor));
|
||||
}
|
||||
|
||||
@@ -5,35 +5,55 @@ using Newtonsoft.Json;
|
||||
|
||||
namespace CodexContractsPlugin.Marketplace
|
||||
{
|
||||
public partial class Request : RequestBase
|
||||
public interface IHasBlock
|
||||
{
|
||||
BlockTimeEntry Block { get; set; }
|
||||
}
|
||||
|
||||
public partial class Request : RequestBase, IHasBlock
|
||||
{
|
||||
[JsonIgnore]
|
||||
public BlockTimeEntry Block { get; set; }
|
||||
public byte[] RequestId { get; set; }
|
||||
|
||||
public EthAddress ClientAddress { get { return new EthAddress(Client); } }
|
||||
|
||||
[JsonIgnore]
|
||||
public string Id
|
||||
{
|
||||
get
|
||||
{
|
||||
return BitConverter.ToString(RequestId).Replace("-", "").ToLowerInvariant();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
public partial class RequestFulfilledEventDTO
|
||||
public partial class RequestFulfilledEventDTO : IHasBlock
|
||||
{
|
||||
[JsonIgnore]
|
||||
public BlockTimeEntry Block { get; set; }
|
||||
}
|
||||
|
||||
public partial class RequestCancelledEventDTO
|
||||
public partial class RequestCancelledEventDTO : IHasBlock
|
||||
{
|
||||
[JsonIgnore]
|
||||
public BlockTimeEntry Block { get; set; }
|
||||
}
|
||||
|
||||
public partial class SlotFilledEventDTO
|
||||
public partial class RequestFailedEventDTO : IHasBlock
|
||||
{
|
||||
[JsonIgnore]
|
||||
public BlockTimeEntry Block { get; set; }
|
||||
}
|
||||
|
||||
public partial class SlotFilledEventDTO : IHasBlock
|
||||
{
|
||||
[JsonIgnore]
|
||||
public BlockTimeEntry Block { get; set; }
|
||||
public EthAddress Host { get; set; }
|
||||
}
|
||||
|
||||
public partial class SlotFreedEventDTO
|
||||
public partial class SlotFreedEventDTO : IHasBlock
|
||||
{
|
||||
[JsonIgnore]
|
||||
public BlockTimeEntry Block { get; set; }
|
||||
|
||||
File diff suppressed because one or more lines are too long
@@ -0,0 +1,108 @@
|
||||
namespace CodexContractsPlugin
|
||||
{
|
||||
public class SelfUpdater
|
||||
{
|
||||
public void Update(string abi, string bytecode)
|
||||
{
|
||||
var filePath = GetMarketplaceFilePath();
|
||||
var content = GenerateContent(abi, bytecode);
|
||||
var contentLines = content.Split("\r\n");
|
||||
|
||||
var beginWith = new string[]
|
||||
{
|
||||
"using Nethereum.ABI.FunctionEncoding.Attributes;",
|
||||
"using Nethereum.Contracts;",
|
||||
"using System.Numerics;",
|
||||
"",
|
||||
"// Generated code, do not modify.",
|
||||
"",
|
||||
"#pragma warning disable CS8618 // Non-nullable field must contain a non-null value when exiting constructor. Consider declaring as nullable.",
|
||||
"namespace CodexContractsPlugin.Marketplace",
|
||||
"{"
|
||||
};
|
||||
|
||||
var endWith = new string[]
|
||||
{
|
||||
"}",
|
||||
"#pragma warning restore CS8618 // Non-nullable field must contain a non-null value when exiting constructor. Consider declaring as nullable."
|
||||
};
|
||||
|
||||
File.Delete(filePath);
|
||||
File.WriteAllLines(filePath,
|
||||
beginWith.Concat(
|
||||
contentLines.Concat(
|
||||
endWith))
|
||||
);
|
||||
|
||||
throw new Exception("Oh no! CodexContracts were updated. Current build of CodexContractsPlugin is incompatible. " +
|
||||
"But fear not! SelfUpdater.cs has automatically updated the plugin. Just rebuild and rerun and it should work. " +
|
||||
"Just in case, manual update instructions are found here: 'CodexContractsPlugin/Marketplace/README.md'.");
|
||||
}
|
||||
|
||||
private string GetMarketplaceFilePath()
|
||||
{
|
||||
var here = Directory.GetCurrentDirectory();
|
||||
while (true)
|
||||
{
|
||||
var path = GetMarketplaceFile(here);
|
||||
if (path != null) return path;
|
||||
|
||||
var parent = Directory.GetParent(here);
|
||||
var up = parent?.FullName;
|
||||
if (up == null || up == here) throw new Exception("Unable to locate ProjectPlugins folder. Unable to update contracts.");
|
||||
here = up;
|
||||
}
|
||||
}
|
||||
|
||||
private string? GetMarketplaceFile(string root)
|
||||
{
|
||||
var path = Path.Combine(root, "ProjectPlugins", "CodexContractsPlugin", "Marketplace", "Marketplace.cs");
|
||||
if (File.Exists(path)) return path;
|
||||
return null;
|
||||
}
|
||||
|
||||
private string GenerateContent(string abi, string bytecode)
|
||||
{
|
||||
var deserializer = new Nethereum.Generators.Net.GeneratorModelABIDeserialiser();
|
||||
var abiModel = deserializer.DeserialiseABI(abi);
|
||||
var abiCtor = abiModel.Constructor;
|
||||
var c = new Nethereum.Generators.CQS.ContractDeploymentCQSMessageGenerator(abiCtor, "namespace", bytecode, "Marketplace", Nethereum.Generators.Core.CodeGenLanguage.CSharp);
|
||||
var lines = "";
|
||||
lines += c.GenerateClass();
|
||||
lines += "\r\n";
|
||||
|
||||
foreach (var eventAbi in abiModel.Events)
|
||||
{
|
||||
var d = new Nethereum.Generators.DTOs.EventDTOGenerator(eventAbi, "namespace", Nethereum.Generators.Core.CodeGenLanguage.CSharp);
|
||||
lines += d.GenerateClass();
|
||||
lines += "\r\n";
|
||||
}
|
||||
|
||||
foreach (var errorAbi in abiModel.Errors)
|
||||
{
|
||||
var e = new Nethereum.Generators.DTOs.ErrorDTOGenerator(errorAbi, "namespace", Nethereum.Generators.Core.CodeGenLanguage.CSharp);
|
||||
lines += e.GenerateClass();
|
||||
lines += "\r\n";
|
||||
}
|
||||
|
||||
foreach (var funcAbi in abiModel.Functions)
|
||||
{
|
||||
var f = new Nethereum.Generators.DTOs.FunctionOutputDTOGenerator(funcAbi, "namespace", Nethereum.Generators.Core.CodeGenLanguage.CSharp);
|
||||
var ff = new Nethereum.Generators.CQS.FunctionCQSMessageGenerator(funcAbi, "namespace", "funcoutput", Nethereum.Generators.Core.CodeGenLanguage.CSharp);
|
||||
lines += f.GenerateClass();
|
||||
lines += "\r\n";
|
||||
lines += ff.GenerateClass();
|
||||
lines += "\r\n";
|
||||
}
|
||||
|
||||
foreach (var structAbi in abiModel.Structs)
|
||||
{
|
||||
var g = new Nethereum.Generators.DTOs.StructTypeGenerator(structAbi, "namespace", Nethereum.Generators.Core.CodeGenLanguage.CSharp);
|
||||
lines += g.GenerateClass();
|
||||
lines += "\r\n";
|
||||
}
|
||||
|
||||
return lines;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,11 +1,13 @@
|
||||
using Core;
|
||||
using KubernetesWorkflow;
|
||||
using KubernetesWorkflow.Types;
|
||||
using Utils;
|
||||
|
||||
namespace CodexDiscordBotPlugin
|
||||
{
|
||||
public class CodexDiscordBotPlugin : IProjectPlugin, IHasLogPrefix, IHasMetadata
|
||||
{
|
||||
private const string ExpectedStartupMessage = "Debug option is set. Discord connection disabled!";
|
||||
private readonly IPluginTools tools;
|
||||
|
||||
public CodexDiscordBotPlugin(IPluginTools tools)
|
||||
@@ -46,14 +48,59 @@ namespace CodexDiscordBotPlugin
|
||||
var startupConfig = new StartupConfig();
|
||||
startupConfig.NameOverride = config.Name;
|
||||
startupConfig.Add(config);
|
||||
return workflow.Start(1, new DiscordBotContainerRecipe(), startupConfig).WaitForOnline();
|
||||
var pod = workflow.Start(1, new DiscordBotContainerRecipe(), startupConfig).WaitForOnline();
|
||||
WaitForStartupMessage(workflow, pod);
|
||||
workflow.CreateCrashWatcher(pod.Containers.Single()).Start();
|
||||
return pod;
|
||||
}
|
||||
|
||||
private RunningPod StartRewarderContainer(IStartupWorkflow workflow, RewarderBotStartupConfig config)
|
||||
{
|
||||
var startupConfig = new StartupConfig();
|
||||
startupConfig.NameOverride = config.Name;
|
||||
startupConfig.Add(config);
|
||||
return workflow.Start(1, new RewarderBotContainerRecipe(), startupConfig).WaitForOnline();
|
||||
var pod = workflow.Start(1, new RewarderBotContainerRecipe(), startupConfig).WaitForOnline();
|
||||
workflow.CreateCrashWatcher(pod.Containers.Single()).Start();
|
||||
return pod;
|
||||
}
|
||||
|
||||
private void WaitForStartupMessage(IStartupWorkflow workflow, RunningPod pod)
|
||||
{
|
||||
var finder = new LogLineFinder(ExpectedStartupMessage, workflow);
|
||||
Time.WaitUntil(() =>
|
||||
{
|
||||
finder.FindLine(pod);
|
||||
return finder.Found;
|
||||
}, nameof(WaitForStartupMessage));
|
||||
}
|
||||
|
||||
public class LogLineFinder : LogHandler
|
||||
{
|
||||
private readonly string message;
|
||||
private readonly IStartupWorkflow workflow;
|
||||
|
||||
public LogLineFinder(string message, IStartupWorkflow workflow)
|
||||
{
|
||||
this.message = message;
|
||||
this.workflow = workflow;
|
||||
}
|
||||
|
||||
public void FindLine(RunningPod pod)
|
||||
{
|
||||
Found = false;
|
||||
foreach (var c in pod.Containers)
|
||||
{
|
||||
workflow.DownloadContainerLog(c, this);
|
||||
if (Found) return;
|
||||
}
|
||||
}
|
||||
|
||||
public bool Found { get; private set; }
|
||||
|
||||
protected override void ProcessLine(string line)
|
||||
{
|
||||
if (!Found && line.Contains(message)) Found = true;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -7,7 +7,7 @@ namespace CodexDiscordBotPlugin
|
||||
public class DiscordBotContainerRecipe : ContainerRecipeFactory
|
||||
{
|
||||
public override string AppName => "discordbot-bibliotech";
|
||||
public override string Image => "codexstorage/codex-discordbot:sha-8c64352";
|
||||
public override string Image => "codexstorage/codex-discordbot:sha-8033da1";
|
||||
|
||||
public static string RewardsPort = "bot_rewards_port";
|
||||
|
||||
@@ -33,6 +33,8 @@ namespace CodexDiscordBotPlugin
|
||||
AddEnvVar("CODEXCONTRACTS_TOKENADDRESS", gethInfo.TokenAddress);
|
||||
AddEnvVar("CODEXCONTRACTS_ABI", gethInfo.Abi);
|
||||
|
||||
AddEnvVar("NODISCORD", "1");
|
||||
|
||||
AddInternalPortAndVar("REWARDAPIPORT", RewardsPort);
|
||||
|
||||
if (!string.IsNullOrEmpty(config.DataPath))
|
||||
|
||||
@@ -27,8 +27,9 @@
|
||||
|
||||
public class RewarderBotStartupConfig
|
||||
{
|
||||
public RewarderBotStartupConfig(string discordBotHost, int discordBotPort, string intervalMinutes, DateTime historyStartUtc, DiscordBotGethInfo gethInfo, string? dataPath)
|
||||
public RewarderBotStartupConfig(string name, string discordBotHost, int discordBotPort, int intervalMinutes, DateTime historyStartUtc, DiscordBotGethInfo gethInfo, string? dataPath)
|
||||
{
|
||||
Name = name;
|
||||
DiscordBotHost = discordBotHost;
|
||||
DiscordBotPort = discordBotPort;
|
||||
IntervalMinutes = intervalMinutes;
|
||||
@@ -37,9 +38,10 @@
|
||||
DataPath = dataPath;
|
||||
}
|
||||
|
||||
public string Name { get; }
|
||||
public string DiscordBotHost { get; }
|
||||
public int DiscordBotPort { get; }
|
||||
public string IntervalMinutes { get; }
|
||||
public int IntervalMinutes { get; }
|
||||
public DateTime HistoryStartUtc { get; }
|
||||
public DiscordBotGethInfo GethInfo { get; }
|
||||
public string? DataPath { get; set; }
|
||||
|
||||
@@ -7,7 +7,7 @@ namespace CodexDiscordBotPlugin
|
||||
public class RewarderBotContainerRecipe : ContainerRecipeFactory
|
||||
{
|
||||
public override string AppName => "discordbot-rewarder";
|
||||
public override string Image => "codexstorage/codex-rewarderbot:sha-2ab84e2";
|
||||
public override string Image => "codexstorage/codex-rewarderbot:sha-8033da1";
|
||||
|
||||
protected override void Initialize(StartupConfig startupConfig)
|
||||
{
|
||||
@@ -17,7 +17,7 @@ namespace CodexDiscordBotPlugin
|
||||
|
||||
AddEnvVar("DISCORDBOTHOST", config.DiscordBotHost);
|
||||
AddEnvVar("DISCORDBOTPORT", config.DiscordBotPort.ToString());
|
||||
AddEnvVar("INTERVALMINUTES", config.IntervalMinutes);
|
||||
AddEnvVar("INTERVALMINUTES", config.IntervalMinutes.ToString());
|
||||
var offset = new DateTimeOffset(config.HistoryStartUtc);
|
||||
AddEnvVar("CHECKHISTORY", offset.ToUnixTimeSeconds().ToString());
|
||||
|
||||
|
||||
@@ -132,19 +132,6 @@ namespace CodexPlugin
|
||||
return workflow.GetPodInfo(Container);
|
||||
}
|
||||
|
||||
public void LogDiskSpace(string msg)
|
||||
{
|
||||
try
|
||||
{
|
||||
var diskInfo = tools.CreateWorkflow().ExecuteCommand(Container.Containers.Single(), "df", "--sync");
|
||||
Log($"{msg} - Disk info: {diskInfo}");
|
||||
}
|
||||
catch (Exception e)
|
||||
{
|
||||
Log("Failed to get disk info: " + e);
|
||||
}
|
||||
}
|
||||
|
||||
public void DeleteRepoFolder()
|
||||
{
|
||||
try
|
||||
@@ -163,13 +150,13 @@ namespace CodexPlugin
|
||||
|
||||
private T OnCodex<T>(Func<CodexApi, Task<T>> action)
|
||||
{
|
||||
var result = tools.CreateHttp(CheckContainerCrashed).OnClient(client => CallCodex(client, action));
|
||||
var result = tools.CreateHttp(GetHttpId(), CheckContainerCrashed).OnClient(client => CallCodex(client, action));
|
||||
return result;
|
||||
}
|
||||
|
||||
private T OnCodex<T>(Func<CodexApi, Task<T>> action, Retry retry)
|
||||
{
|
||||
var result = tools.CreateHttp(CheckContainerCrashed).OnClient(client => CallCodex(client, action), retry);
|
||||
var result = tools.CreateHttp(GetHttpId(), CheckContainerCrashed).OnClient(client => CallCodex(client, action), retry);
|
||||
return result;
|
||||
}
|
||||
|
||||
@@ -196,7 +183,7 @@ namespace CodexPlugin
|
||||
private IEndpoint GetEndpoint()
|
||||
{
|
||||
return tools
|
||||
.CreateHttp(CheckContainerCrashed)
|
||||
.CreateHttp(GetHttpId(), CheckContainerCrashed)
|
||||
.CreateEndpoint(GetAddress(), "/api/codex/v1/", Container.Name);
|
||||
}
|
||||
|
||||
@@ -205,6 +192,11 @@ namespace CodexPlugin
|
||||
return Container.Containers.Single().GetAddress(log, CodexContainerRecipe.ApiPortTag);
|
||||
}
|
||||
|
||||
private string GetHttpId()
|
||||
{
|
||||
return GetAddress().ToString();
|
||||
}
|
||||
|
||||
private void CheckContainerCrashed(HttpClient client)
|
||||
{
|
||||
if (CrashWatcher.HasContainerCrashed()) throw new Exception($"Container {GetName()} has crashed.");
|
||||
@@ -227,8 +219,6 @@ namespace CodexPlugin
|
||||
$"(HTTP timeout = {Time.FormatDuration(timeSet.HttpCallTimeout())}) " +
|
||||
$"Checking if node responds to debug/info...");
|
||||
|
||||
LogDiskSpace("After retry failure");
|
||||
|
||||
try
|
||||
{
|
||||
var debugInfo = GetDebugInfo();
|
||||
|
||||
@@ -7,7 +7,7 @@ namespace CodexPlugin
|
||||
{
|
||||
public class CodexContainerRecipe : ContainerRecipeFactory
|
||||
{
|
||||
private const string DefaultDockerImage = "codexstorage/nim-codex:sha-b89493e-dist-tests";
|
||||
private const string DefaultDockerImage = "codexstorage/nim-codex:sha-1e2ad95-dist-tests";
|
||||
|
||||
public const string ApiPortTag = "codex_api_port";
|
||||
public const string ListenPortTag = "codex_listen_port";
|
||||
@@ -109,8 +109,9 @@ namespace CodexPlugin
|
||||
|
||||
// Custom scripting in the Codex test image will write this variable to a private-key file,
|
||||
// and pass the correct filename to Codex.
|
||||
AddEnvVar("PRIV_KEY", marketplaceSetup.EthAccount.PrivateKey);
|
||||
Additional(marketplaceSetup.EthAccount);
|
||||
var account = marketplaceSetup.EthAccountSetup.GetNew();
|
||||
AddEnvVar("PRIV_KEY", account.PrivateKey);
|
||||
Additional(account);
|
||||
|
||||
SetCommandOverride(marketplaceSetup);
|
||||
if (marketplaceSetup.IsValidator)
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
using Core;
|
||||
using CodexPlugin.Hooks;
|
||||
using Core;
|
||||
using FileUtils;
|
||||
using GethPlugin;
|
||||
using KubernetesWorkflow;
|
||||
@@ -12,6 +13,7 @@ namespace CodexPlugin
|
||||
public interface ICodexNode : IHasContainer, IHasMetricsScrapeTarget, IHasEthAddress
|
||||
{
|
||||
string GetName();
|
||||
string GetPeerId();
|
||||
DebugInfo GetDebugInfo();
|
||||
DebugPeer GetDebugPeer(string peerId);
|
||||
ContentId UploadFile(TrackedFile file);
|
||||
@@ -26,32 +28,49 @@ namespace CodexPlugin
|
||||
CrashWatcher CrashWatcher { get; }
|
||||
PodInfo GetPodInfo();
|
||||
ITransferSpeeds TransferSpeeds { get; }
|
||||
EthAccount EthAccount { get; }
|
||||
|
||||
/// <summary>
|
||||
/// Warning! The node is not usable after this.
|
||||
/// TODO: Replace with delete-blocks debug call once available in Codex.
|
||||
/// </summary>
|
||||
void DeleteRepoFolder();
|
||||
|
||||
void Stop(bool waitTillStopped);
|
||||
}
|
||||
|
||||
public class CodexNode : ICodexNode
|
||||
{
|
||||
private const string UploadFailedMessage = "Unable to store block";
|
||||
private readonly ILog log;
|
||||
private readonly IPluginTools tools;
|
||||
private readonly EthAddress? ethAddress;
|
||||
private readonly ICodexNodeHooks hooks;
|
||||
private readonly EthAccount? ethAccount;
|
||||
private readonly TransferSpeeds transferSpeeds;
|
||||
private string peerId = string.Empty;
|
||||
private string nodeId = string.Empty;
|
||||
|
||||
public CodexNode(IPluginTools tools, CodexAccess codexAccess, CodexNodeGroup group, IMarketplaceAccess marketplaceAccess, EthAddress? ethAddress)
|
||||
public CodexNode(IPluginTools tools, CodexAccess codexAccess, CodexNodeGroup group, IMarketplaceAccess marketplaceAccess, ICodexNodeHooks hooks, EthAccount? ethAccount)
|
||||
{
|
||||
this.tools = tools;
|
||||
this.ethAddress = ethAddress;
|
||||
this.ethAccount = ethAccount;
|
||||
CodexAccess = codexAccess;
|
||||
Group = group;
|
||||
Marketplace = marketplaceAccess;
|
||||
this.hooks = hooks;
|
||||
Version = new DebugInfoVersion();
|
||||
transferSpeeds = new TransferSpeeds();
|
||||
|
||||
log = new LogPrefixer(tools.GetLog(), $"{GetName()} ");
|
||||
}
|
||||
|
||||
public void Awake()
|
||||
{
|
||||
hooks.OnNodeStarting(Container.Recipe.RecipeCreatedUtc, Container.Recipe.Image, ethAccount);
|
||||
}
|
||||
|
||||
public void Initialize()
|
||||
{
|
||||
hooks.OnNodeStarted(peerId, nodeId);
|
||||
}
|
||||
|
||||
public RunningPod Pod { get { return CodexAccess.Container; } }
|
||||
@@ -76,8 +95,17 @@ namespace CodexPlugin
|
||||
{
|
||||
get
|
||||
{
|
||||
if (ethAddress == null) throw new Exception("Marketplace is not enabled for this Codex node. Please start it with the option '.EnableMarketplace(...)' to enable it.");
|
||||
return ethAddress;
|
||||
EnsureMarketplace();
|
||||
return ethAccount!.EthAddress;
|
||||
}
|
||||
}
|
||||
|
||||
public EthAccount EthAccount
|
||||
{
|
||||
get
|
||||
{
|
||||
EnsureMarketplace();
|
||||
return ethAccount!;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -86,6 +114,11 @@ namespace CodexPlugin
|
||||
return Container.Name;
|
||||
}
|
||||
|
||||
public string GetPeerId()
|
||||
{
|
||||
return peerId;
|
||||
}
|
||||
|
||||
public DebugInfo GetDebugInfo()
|
||||
{
|
||||
var debugInfo = CodexAccess.GetDebugInfo();
|
||||
@@ -106,27 +139,29 @@ namespace CodexPlugin
|
||||
|
||||
public ContentId UploadFile(TrackedFile file, Action<Failure> onFailure)
|
||||
{
|
||||
CodexAccess.LogDiskSpace("Before upload");
|
||||
|
||||
using var fileStream = File.OpenRead(file.Filename);
|
||||
var uniqueId = Guid.NewGuid().ToString();
|
||||
var size = file.GetFilesize();
|
||||
|
||||
hooks.OnFileUploading(uniqueId, size);
|
||||
|
||||
var logMessage = $"Uploading file {file.Describe()}...";
|
||||
Log(logMessage);
|
||||
var measurement = Stopwatch.Measure(tools.GetLog(), logMessage, () =>
|
||||
var measurement = Stopwatch.Measure(log, logMessage, () =>
|
||||
{
|
||||
return CodexAccess.UploadFile(fileStream, onFailure);
|
||||
});
|
||||
|
||||
var response = measurement.Value;
|
||||
transferSpeeds.AddUploadSample(file.GetFilesize(), measurement.Duration);
|
||||
transferSpeeds.AddUploadSample(size, measurement.Duration);
|
||||
|
||||
if (string.IsNullOrEmpty(response)) FrameworkAssert.Fail("Received empty response.");
|
||||
if (response.StartsWith(UploadFailedMessage)) FrameworkAssert.Fail("Node failed to store block.");
|
||||
|
||||
Log($"Uploaded file. Received contentId: '{response}'.");
|
||||
CodexAccess.LogDiskSpace("After upload");
|
||||
Log($"Uploaded file {file.Describe()}. Received contentId: '{response}'.");
|
||||
|
||||
return new ContentId(response);
|
||||
var cid = new ContentId(response);
|
||||
hooks.OnFileUploaded(uniqueId, size, cid);
|
||||
return cid;
|
||||
}
|
||||
|
||||
public TrackedFile? DownloadContent(ContentId contentId, string fileLabel = "")
|
||||
@@ -136,12 +171,17 @@ namespace CodexPlugin
|
||||
|
||||
public TrackedFile? DownloadContent(ContentId contentId, Action<Failure> onFailure, string fileLabel = "")
|
||||
{
|
||||
var logMessage = $"Downloading for contentId: '{contentId.Id}'...";
|
||||
Log(logMessage);
|
||||
var file = tools.GetFileManager().CreateEmptyFile(fileLabel);
|
||||
var measurement = Stopwatch.Measure(tools.GetLog(), logMessage, () => DownloadToFile(contentId.Id, file, onFailure));
|
||||
transferSpeeds.AddDownloadSample(file.GetFilesize(), measurement);
|
||||
Log($"Downloaded file {file.Describe()} to '{file.Filename}'.");
|
||||
hooks.OnFileDownloading(contentId);
|
||||
Log($"Downloading '{contentId}'...");
|
||||
|
||||
var logMessage = $"Downloaded '{contentId}' to '{file.Filename}'";
|
||||
var measurement = Stopwatch.Measure(log, logMessage, () => DownloadToFile(contentId.Id, file, onFailure));
|
||||
|
||||
var size = file.GetFilesize();
|
||||
transferSpeeds.AddDownloadSample(size, measurement);
|
||||
hooks.OnFileDownloaded(size, contentId);
|
||||
|
||||
return file;
|
||||
}
|
||||
|
||||
@@ -179,6 +219,8 @@ namespace CodexPlugin
|
||||
public void Stop(bool waitTillStopped)
|
||||
{
|
||||
Log("Stopping...");
|
||||
hooks.OnNodeStopping();
|
||||
|
||||
CrashWatcher.Stop();
|
||||
Group.Stop(this, waitTillStopped);
|
||||
}
|
||||
@@ -186,7 +228,8 @@ namespace CodexPlugin
|
||||
public void EnsureOnlineGetVersionResponse()
|
||||
{
|
||||
var debugInfo = Time.Retry(CodexAccess.GetDebugInfo, "ensure online");
|
||||
var nodePeerId = debugInfo.Id;
|
||||
peerId = debugInfo.Id;
|
||||
nodeId = debugInfo.Table.LocalNode.NodeId;
|
||||
var nodeName = CodexAccess.Container.Name;
|
||||
|
||||
if (!debugInfo.Version.IsValid())
|
||||
@@ -194,9 +237,10 @@ namespace CodexPlugin
|
||||
throw new Exception($"Invalid version information received from Codex node {GetName()}: {debugInfo.Version}");
|
||||
}
|
||||
|
||||
var log = tools.GetLog();
|
||||
log.AddStringReplace(nodePeerId, nodeName);
|
||||
log.AddStringReplace(peerId, nodeName);
|
||||
log.AddStringReplace(CodexUtils.ToShortId(peerId), nodeName);
|
||||
log.AddStringReplace(debugInfo.Table.LocalNode.NodeId, nodeName);
|
||||
log.AddStringReplace(CodexUtils.ToShortId(debugInfo.Table.LocalNode.NodeId), nodeName);
|
||||
Version = debugInfo.Version;
|
||||
}
|
||||
|
||||
@@ -214,8 +258,6 @@ namespace CodexPlugin
|
||||
|
||||
private void DownloadToFile(string contentId, TrackedFile file, Action<Failure> onFailure)
|
||||
{
|
||||
CodexAccess.LogDiskSpace("Before download");
|
||||
|
||||
using var fileStream = File.OpenWrite(file.Filename);
|
||||
try
|
||||
{
|
||||
@@ -227,13 +269,16 @@ namespace CodexPlugin
|
||||
Log($"Failed to download file '{contentId}'.");
|
||||
throw;
|
||||
}
|
||||
}
|
||||
|
||||
CodexAccess.LogDiskSpace("After download");
|
||||
private void EnsureMarketplace()
|
||||
{
|
||||
if (ethAccount == null) throw new Exception("Marketplace is not enabled for this Codex node. Please start it with the option '.EnableMarketplace(...)' to enable it.");
|
||||
}
|
||||
|
||||
private void Log(string msg)
|
||||
{
|
||||
tools.GetLog().Log($"{GetName()}: {msg}");
|
||||
log.Log(msg);
|
||||
}
|
||||
|
||||
private void DoNothing(Failure failure)
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
using Core;
|
||||
using CodexPlugin.Hooks;
|
||||
using Core;
|
||||
using GethPlugin;
|
||||
using KubernetesWorkflow;
|
||||
using KubernetesWorkflow.Types;
|
||||
@@ -14,30 +15,34 @@ namespace CodexPlugin
|
||||
public class CodexNodeFactory : ICodexNodeFactory
|
||||
{
|
||||
private readonly IPluginTools tools;
|
||||
private readonly CodexHooksFactory codexHooksFactory;
|
||||
|
||||
public CodexNodeFactory(IPluginTools tools)
|
||||
public CodexNodeFactory(IPluginTools tools, CodexHooksFactory codexHooksFactory)
|
||||
{
|
||||
this.tools = tools;
|
||||
this.codexHooksFactory = codexHooksFactory;
|
||||
}
|
||||
|
||||
public CodexNode CreateOnlineCodexNode(CodexAccess access, CodexNodeGroup group)
|
||||
{
|
||||
var ethAddress = GetEthAddress(access);
|
||||
var marketplaceAccess = GetMarketplaceAccess(access, ethAddress);
|
||||
return new CodexNode(tools, access, group, marketplaceAccess, ethAddress);
|
||||
var ethAccount = GetEthAccount(access);
|
||||
var hooks = codexHooksFactory.CreateHooks(access.Container.Name);
|
||||
|
||||
var marketplaceAccess = GetMarketplaceAccess(access, ethAccount, hooks);
|
||||
return new CodexNode(tools, access, group, marketplaceAccess, hooks, ethAccount);
|
||||
}
|
||||
|
||||
private IMarketplaceAccess GetMarketplaceAccess(CodexAccess codexAccess, EthAddress? ethAddress)
|
||||
private IMarketplaceAccess GetMarketplaceAccess(CodexAccess codexAccess, EthAccount? ethAccount, ICodexNodeHooks hooks)
|
||||
{
|
||||
if (ethAddress == null) return new MarketplaceUnavailable();
|
||||
return new MarketplaceAccess(tools.GetLog(), codexAccess);
|
||||
if (ethAccount == null) return new MarketplaceUnavailable();
|
||||
return new MarketplaceAccess(tools.GetLog(), codexAccess, hooks);
|
||||
}
|
||||
|
||||
private EthAddress? GetEthAddress(CodexAccess access)
|
||||
private EthAccount? GetEthAccount(CodexAccess access)
|
||||
{
|
||||
var ethAccount = access.Container.Containers.Single().Recipe.Additionals.Get<EthAccount>();
|
||||
if (ethAccount == null) return null;
|
||||
return ethAccount.EthAddress;
|
||||
return ethAccount;
|
||||
}
|
||||
|
||||
public CrashWatcher CreateCrashWatcher(RunningContainer c)
|
||||
|
||||
@@ -79,13 +79,16 @@ namespace CodexPlugin
|
||||
}
|
||||
|
||||
Version = first;
|
||||
foreach (var node in Nodes) node.Initialize();
|
||||
}
|
||||
|
||||
private CodexNode CreateOnlineCodexNode(RunningPod c, IPluginTools tools, ICodexNodeFactory factory)
|
||||
{
|
||||
var watcher = factory.CreateCrashWatcher(c.Containers.Single());
|
||||
var access = new CodexAccess(tools, c, watcher);
|
||||
return factory.CreateOnlineCodexNode(access, this);
|
||||
var node = factory.CreateOnlineCodexNode(access, this);
|
||||
node.Awake();
|
||||
return node;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,3 +1,4 @@
|
||||
using CodexPlugin.Hooks;
|
||||
using Core;
|
||||
using KubernetesWorkflow.Types;
|
||||
|
||||
@@ -57,6 +58,11 @@ namespace CodexPlugin
|
||||
}
|
||||
}
|
||||
|
||||
public void SetCodexHooksProvider(ICodexHooksProvider hooksProvider)
|
||||
{
|
||||
codexStarter.HooksFactory.Provider = hooksProvider;
|
||||
}
|
||||
|
||||
private CodexSetup GetSetup(int numberOfNodes, Action<ICodexSetup> setup)
|
||||
{
|
||||
var codexSetup = new CodexSetup(numberOfNodes);
|
||||
|
||||
@@ -29,6 +29,7 @@
|
||||
<ItemGroup>
|
||||
<ProjectReference Include="..\..\Framework\Core\Core.csproj" />
|
||||
<ProjectReference Include="..\..\Framework\KubernetesWorkflow\KubernetesWorkflow.csproj" />
|
||||
<ProjectReference Include="..\..\Framework\OverwatchTranscript\OverwatchTranscript.csproj" />
|
||||
<ProjectReference Include="..\CodexContractsPlugin\CodexContractsPlugin.csproj" />
|
||||
<ProjectReference Include="..\GethPlugin\GethPlugin.csproj" />
|
||||
<ProjectReference Include="..\MetricsPlugin\MetricsPlugin.csproj" />
|
||||
|
||||
@@ -169,7 +169,7 @@ namespace CodexPlugin
|
||||
public bool IsValidator { get; private set; }
|
||||
public Ether InitialEth { get; private set; } = 0.Eth();
|
||||
public TestToken InitialTestTokens { get; private set; } = 0.Tst();
|
||||
public EthAccount EthAccount { get; private set; } = EthAccount.GenerateNew();
|
||||
public EthAccountSetup EthAccountSetup { get; } = new EthAccountSetup();
|
||||
|
||||
public IMarketplaceSetup AsStorageNode()
|
||||
{
|
||||
@@ -185,7 +185,7 @@ namespace CodexPlugin
|
||||
|
||||
public IMarketplaceSetup WithAccount(EthAccount account)
|
||||
{
|
||||
EthAccount = account;
|
||||
EthAccountSetup.Pin(account);
|
||||
return this;
|
||||
}
|
||||
|
||||
@@ -201,10 +201,42 @@ namespace CodexPlugin
|
||||
var result = "[(clientNode)"; // When marketplace is enabled, being a clientNode is implicit.
|
||||
result += IsStorageNode ? "(storageNode)" : "()";
|
||||
result += IsValidator ? "(validator)" : "() ";
|
||||
result += $"Address: '{EthAccount.EthAddress}' ";
|
||||
result += $"{InitialEth.Eth} / {InitialTestTokens}";
|
||||
result += $"Pinned address: '{EthAccountSetup}' ";
|
||||
result += $"{InitialEth} / {InitialTestTokens}";
|
||||
result += "] ";
|
||||
return result;
|
||||
}
|
||||
}
|
||||
|
||||
public class EthAccountSetup
|
||||
{
|
||||
private readonly List<EthAccount> accounts = new List<EthAccount>();
|
||||
private bool pinned = false;
|
||||
|
||||
public void Pin(EthAccount account)
|
||||
{
|
||||
accounts.Add(account);
|
||||
pinned = true;
|
||||
}
|
||||
|
||||
public EthAccount GetNew()
|
||||
{
|
||||
if (pinned) return accounts.Last();
|
||||
|
||||
var a = EthAccount.GenerateNew();
|
||||
accounts.Add(a);
|
||||
return a;
|
||||
}
|
||||
|
||||
public EthAccount[] GetAll()
|
||||
{
|
||||
return accounts.ToArray();
|
||||
}
|
||||
|
||||
public override string ToString()
|
||||
{
|
||||
if (!accounts.Any()) return "NoEthAccounts";
|
||||
return string.Join(",", accounts.Select(a => a.ToString()).ToArray());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,4 +1,6 @@
|
||||
using Core;
|
||||
using CodexPlugin.Hooks;
|
||||
using Core;
|
||||
using GethPlugin;
|
||||
using KubernetesWorkflow;
|
||||
using KubernetesWorkflow.Types;
|
||||
using Logging;
|
||||
@@ -19,6 +21,8 @@ namespace CodexPlugin
|
||||
apiChecker = new ApiChecker(pluginTools);
|
||||
}
|
||||
|
||||
public CodexHooksFactory HooksFactory { get; } = new CodexHooksFactory();
|
||||
|
||||
public RunningPod[] BringOnline(CodexSetup codexSetup)
|
||||
{
|
||||
LogSeparator();
|
||||
@@ -34,7 +38,8 @@ namespace CodexPlugin
|
||||
{
|
||||
var podInfo = GetPodInfo(rc);
|
||||
var podInfos = string.Join(", ", rc.Containers.Select(c => $"Container: '{c.Name}' PodLabel: '{c.RunningPod.StartResult.Deployment.PodLabel}' runs at '{podInfo.K8SNodeName}'={podInfo.Ip}"));
|
||||
Log($"Started {codexSetup.NumberOfNodes} nodes of image '{containers.First().Containers.First().Recipe.Image}'. ({podInfos})");
|
||||
Log($"Started node with image '{containers.First().Containers.First().Recipe.Image}'. ({podInfos})");
|
||||
LogEthAddress(rc);
|
||||
}
|
||||
LogSeparator();
|
||||
|
||||
@@ -43,7 +48,7 @@ namespace CodexPlugin
|
||||
|
||||
public ICodexNodeGroup WrapCodexContainers(CoreInterface coreInterface, RunningPod[] containers)
|
||||
{
|
||||
var codexNodeFactory = new CodexNodeFactory(pluginTools);
|
||||
var codexNodeFactory = new CodexNodeFactory(pluginTools, HooksFactory);
|
||||
|
||||
var group = CreateCodexGroup(coreInterface, containers, codexNodeFactory);
|
||||
|
||||
@@ -141,6 +146,13 @@ namespace CodexPlugin
|
||||
Log("----------------------------------------------------------------------------");
|
||||
}
|
||||
|
||||
private void LogEthAddress(RunningPod rc)
|
||||
{
|
||||
var account = rc.Containers.First().Recipe.Additionals.Get<EthAccount>();
|
||||
if (account == null) return;
|
||||
Log($"{rc.Name} = {account}");
|
||||
}
|
||||
|
||||
private void Log(string message)
|
||||
{
|
||||
pluginTools.GetLog().Log(message);
|
||||
|
||||
@@ -8,7 +8,7 @@ namespace CodexPlugin
|
||||
public string? NameOverride { get; set; }
|
||||
public ILocation Location { get; set; } = KnownLocations.UnspecifiedLocation;
|
||||
public CodexLogLevel LogLevel { get; set; }
|
||||
public CodexLogCustomTopics? CustomTopics { get; set; } = new CodexLogCustomTopics(CodexLogLevel.Warn, CodexLogLevel.Warn);
|
||||
public CodexLogCustomTopics? CustomTopics { get; set; } = new CodexLogCustomTopics(CodexLogLevel.Info, CodexLogLevel.Warn);
|
||||
public ByteSize? StorageQuota { get; set; }
|
||||
public bool MetricsEnabled { get; set; }
|
||||
public MarketplaceInitialConfig? MarketplaceConfig { get; set; }
|
||||
@@ -50,10 +50,10 @@ namespace CodexPlugin
|
||||
"secure",
|
||||
"chronosstream",
|
||||
"connection",
|
||||
"connmanager",
|
||||
// Removed: "connmanager", is used for transcript peer-dropped event.
|
||||
"websock",
|
||||
"ws-session",
|
||||
"dialer",
|
||||
// Removed: "dialer", is used for transcript successful-dial event.
|
||||
"muxedupgrade",
|
||||
"upgrade",
|
||||
"identify"
|
||||
|
||||
@@ -95,6 +95,11 @@ namespace CodexPlugin
|
||||
|
||||
public string Id { get; }
|
||||
|
||||
public override string ToString()
|
||||
{
|
||||
return Id;
|
||||
}
|
||||
|
||||
public override bool Equals(object? obj)
|
||||
{
|
||||
return obj is ContentId id && Id == id.Id;
|
||||
|
||||
@@ -0,0 +1,24 @@
|
||||
namespace CodexPlugin
|
||||
{
|
||||
public static class CodexUtils
|
||||
{
|
||||
public static string ToShortId(string id)
|
||||
{
|
||||
if (id.Length > 10)
|
||||
{
|
||||
return $"{id[..3]}*{id[^6..]}";
|
||||
}
|
||||
return id;
|
||||
}
|
||||
|
||||
// after update of codex-dht, shortID should be consistent!
|
||||
public static string ToNodeIdShortId(string id)
|
||||
{
|
||||
if (id.Length > 10)
|
||||
{
|
||||
return $"{id[..2]}*{id[^6..]}";
|
||||
}
|
||||
return id;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,4 +1,5 @@
|
||||
using Core;
|
||||
using CodexPlugin.Hooks;
|
||||
using Core;
|
||||
using KubernetesWorkflow.Types;
|
||||
|
||||
namespace CodexPlugin
|
||||
@@ -38,6 +39,11 @@ namespace CodexPlugin
|
||||
return ci.StartCodexNodes(number, s => { });
|
||||
}
|
||||
|
||||
public static void SetCodexHooksProvider(this CoreInterface ci, ICodexHooksProvider hooksProvider)
|
||||
{
|
||||
Plugin(ci).SetCodexHooksProvider(hooksProvider);
|
||||
}
|
||||
|
||||
private static CodexPlugin Plugin(CoreInterface ci)
|
||||
{
|
||||
return ci.GetPlugin<CodexPlugin>();
|
||||
|
||||
@@ -0,0 +1,71 @@
|
||||
using GethPlugin;
|
||||
using Utils;
|
||||
|
||||
namespace CodexPlugin.Hooks
|
||||
{
|
||||
public interface ICodexHooksProvider
|
||||
{
|
||||
ICodexNodeHooks CreateHooks(string nodeName);
|
||||
}
|
||||
|
||||
public class CodexHooksFactory
|
||||
{
|
||||
public ICodexHooksProvider Provider { get; set; } = new DoNothingHooksProvider();
|
||||
|
||||
public ICodexNodeHooks CreateHooks(string nodeName)
|
||||
{
|
||||
return Provider.CreateHooks(nodeName);
|
||||
}
|
||||
}
|
||||
|
||||
public class DoNothingHooksProvider : ICodexHooksProvider
|
||||
{
|
||||
public ICodexNodeHooks CreateHooks(string nodeName)
|
||||
{
|
||||
return new DoNothingCodexHooks();
|
||||
}
|
||||
}
|
||||
|
||||
public class DoNothingCodexHooks : ICodexNodeHooks
|
||||
{
|
||||
public void OnFileDownloaded(ByteSize size, ContentId cid)
|
||||
{
|
||||
}
|
||||
|
||||
public void OnFileDownloading(ContentId cid)
|
||||
{
|
||||
}
|
||||
|
||||
public void OnFileUploaded(string uid, ByteSize size, ContentId cid)
|
||||
{
|
||||
}
|
||||
|
||||
public void OnFileUploading(string uid, ByteSize size)
|
||||
{
|
||||
}
|
||||
|
||||
public void OnNodeStarted(string peerId, string nodeId)
|
||||
{
|
||||
}
|
||||
|
||||
public void OnNodeStarting(DateTime startUtc, string image, EthAccount? ethAccount)
|
||||
{
|
||||
}
|
||||
|
||||
public void OnNodeStopping()
|
||||
{
|
||||
}
|
||||
|
||||
public void OnStorageAvailabilityCreated(StorageAvailability response)
|
||||
{
|
||||
}
|
||||
|
||||
public void OnStorageContractSubmitted(StoragePurchaseContract storagePurchaseContract)
|
||||
{
|
||||
}
|
||||
|
||||
public void OnStorageContractUpdated(StoragePurchase purchaseStatus)
|
||||
{
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,19 @@
|
||||
using GethPlugin;
|
||||
using Utils;
|
||||
|
||||
namespace CodexPlugin.Hooks
|
||||
{
|
||||
public interface ICodexNodeHooks
|
||||
{
|
||||
void OnNodeStarting(DateTime startUtc, string image, EthAccount? ethAccount);
|
||||
void OnNodeStarted(string peerId, string nodeId);
|
||||
void OnNodeStopping();
|
||||
void OnFileUploading(string uid, ByteSize size);
|
||||
void OnFileUploaded(string uid, ByteSize size, ContentId cid);
|
||||
void OnFileDownloading(ContentId cid);
|
||||
void OnFileDownloaded(ByteSize size, ContentId cid);
|
||||
void OnStorageContractSubmitted(StoragePurchaseContract storagePurchaseContract);
|
||||
void OnStorageContractUpdated(StoragePurchase purchaseStatus);
|
||||
void OnStorageAvailabilityCreated(StorageAvailability response);
|
||||
}
|
||||
}
|
||||
@@ -1,5 +1,5 @@
|
||||
using Logging;
|
||||
using Newtonsoft.Json;
|
||||
using CodexPlugin.Hooks;
|
||||
using Logging;
|
||||
using Utils;
|
||||
|
||||
namespace CodexPlugin
|
||||
@@ -7,27 +7,30 @@ namespace CodexPlugin
|
||||
public interface IMarketplaceAccess
|
||||
{
|
||||
string MakeStorageAvailable(StorageAvailability availability);
|
||||
StoragePurchaseContract RequestStorage(StoragePurchaseRequest purchase);
|
||||
IStoragePurchaseContract RequestStorage(StoragePurchaseRequest purchase);
|
||||
}
|
||||
|
||||
public class MarketplaceAccess : IMarketplaceAccess
|
||||
{
|
||||
private readonly ILog log;
|
||||
private readonly CodexAccess codexAccess;
|
||||
private readonly ICodexNodeHooks hooks;
|
||||
|
||||
public MarketplaceAccess(ILog log, CodexAccess codexAccess)
|
||||
public MarketplaceAccess(ILog log, CodexAccess codexAccess, ICodexNodeHooks hooks)
|
||||
{
|
||||
this.log = log;
|
||||
this.codexAccess = codexAccess;
|
||||
this.hooks = hooks;
|
||||
}
|
||||
|
||||
public StoragePurchaseContract RequestStorage(StoragePurchaseRequest purchase)
|
||||
public IStoragePurchaseContract RequestStorage(StoragePurchaseRequest purchase)
|
||||
{
|
||||
purchase.Log(log);
|
||||
|
||||
var response = codexAccess.RequestStorage(purchase);
|
||||
|
||||
if (string.IsNullOrEmpty(response) ||
|
||||
response == "Unable to encode manifest" ||
|
||||
response == "Purchasing not available" ||
|
||||
response == "Expiry required" ||
|
||||
response == "Expiry needs to be in future" ||
|
||||
@@ -38,8 +41,11 @@ namespace CodexPlugin
|
||||
|
||||
Log($"Storage requested successfully. PurchaseId: '{response}'.");
|
||||
|
||||
var contract = new StoragePurchaseContract(log, codexAccess, response, purchase);
|
||||
var contract = new StoragePurchaseContract(log, codexAccess, response, purchase, hooks);
|
||||
contract.WaitForStorageContractSubmitted();
|
||||
|
||||
hooks.OnStorageContractSubmitted(contract);
|
||||
|
||||
return contract;
|
||||
}
|
||||
|
||||
@@ -50,6 +56,7 @@ namespace CodexPlugin
|
||||
var response = codexAccess.SalesAvailability(availability);
|
||||
|
||||
Log($"Storage successfully made available. Id: {response.Id}");
|
||||
hooks.OnStorageAvailabilityCreated(response);
|
||||
|
||||
return response.Id;
|
||||
}
|
||||
@@ -68,7 +75,7 @@ namespace CodexPlugin
|
||||
throw new NotImplementedException();
|
||||
}
|
||||
|
||||
public StoragePurchaseContract RequestStorage(StoragePurchaseRequest purchase)
|
||||
public IStoragePurchaseContract RequestStorage(StoragePurchaseRequest purchase)
|
||||
{
|
||||
Unavailable();
|
||||
throw new NotImplementedException();
|
||||
@@ -80,134 +87,4 @@ namespace CodexPlugin
|
||||
throw new InvalidOperationException();
|
||||
}
|
||||
}
|
||||
|
||||
public class StoragePurchaseContract
|
||||
{
|
||||
private readonly ILog log;
|
||||
private readonly CodexAccess codexAccess;
|
||||
private readonly TimeSpan gracePeriod = TimeSpan.FromSeconds(30);
|
||||
private readonly DateTime contractPendingUtc = DateTime.UtcNow;
|
||||
private DateTime? contractSubmittedUtc = DateTime.UtcNow;
|
||||
private DateTime? contractStartedUtc;
|
||||
private DateTime? contractFinishedUtc;
|
||||
|
||||
public StoragePurchaseContract(ILog log, CodexAccess codexAccess, string purchaseId, StoragePurchaseRequest purchase)
|
||||
{
|
||||
this.log = log;
|
||||
this.codexAccess = codexAccess;
|
||||
PurchaseId = purchaseId;
|
||||
Purchase = purchase;
|
||||
|
||||
ContentId = new ContentId(codexAccess.GetPurchaseStatus(purchaseId).Request.Content.Cid);
|
||||
}
|
||||
|
||||
public string PurchaseId { get; }
|
||||
public StoragePurchaseRequest Purchase { get; }
|
||||
public ContentId ContentId { get; }
|
||||
|
||||
public TimeSpan? PendingToSubmitted => contractSubmittedUtc - contractPendingUtc;
|
||||
public TimeSpan? SubmittedToStarted => contractStartedUtc - contractSubmittedUtc;
|
||||
public TimeSpan? SubmittedToFinished => contractFinishedUtc - contractSubmittedUtc;
|
||||
|
||||
public void WaitForStorageContractSubmitted()
|
||||
{
|
||||
WaitForStorageContractState(gracePeriod, "submitted", sleep: 200);
|
||||
contractSubmittedUtc = DateTime.UtcNow;
|
||||
LogSubmittedDuration();
|
||||
AssertDuration(PendingToSubmitted, gracePeriod, nameof(PendingToSubmitted));
|
||||
}
|
||||
|
||||
public void WaitForStorageContractStarted()
|
||||
{
|
||||
var timeout = Purchase.Expiry + gracePeriod;
|
||||
|
||||
WaitForStorageContractState(timeout, "started");
|
||||
contractStartedUtc = DateTime.UtcNow;
|
||||
LogStartedDuration();
|
||||
AssertDuration(SubmittedToStarted, timeout, nameof(SubmittedToStarted));
|
||||
}
|
||||
|
||||
public void WaitForStorageContractFinished()
|
||||
{
|
||||
if (!contractStartedUtc.HasValue)
|
||||
{
|
||||
WaitForStorageContractStarted();
|
||||
}
|
||||
var currentContractTime = DateTime.UtcNow - contractSubmittedUtc!.Value;
|
||||
var timeout = (Purchase.Duration - currentContractTime) + gracePeriod;
|
||||
WaitForStorageContractState(timeout, "finished");
|
||||
contractFinishedUtc = DateTime.UtcNow;
|
||||
LogFinishedDuration();
|
||||
AssertDuration(SubmittedToFinished, timeout, nameof(SubmittedToFinished));
|
||||
}
|
||||
|
||||
public StoragePurchase GetPurchaseStatus(string purchaseId)
|
||||
{
|
||||
return codexAccess.GetPurchaseStatus(purchaseId);
|
||||
}
|
||||
|
||||
private void WaitForStorageContractState(TimeSpan timeout, string desiredState, int sleep = 1000)
|
||||
{
|
||||
var lastState = "";
|
||||
var waitStart = DateTime.UtcNow;
|
||||
|
||||
Log($"Waiting for {Time.FormatDuration(timeout)} to reach state '{desiredState}'.");
|
||||
while (lastState != desiredState)
|
||||
{
|
||||
var purchaseStatus = codexAccess.GetPurchaseStatus(PurchaseId);
|
||||
var statusJson = JsonConvert.SerializeObject(purchaseStatus);
|
||||
if (purchaseStatus != null && purchaseStatus.State != lastState)
|
||||
{
|
||||
lastState = purchaseStatus.State;
|
||||
log.Debug("Purchase status: " + statusJson);
|
||||
}
|
||||
|
||||
Thread.Sleep(sleep);
|
||||
|
||||
if (lastState == "errored")
|
||||
{
|
||||
FrameworkAssert.Fail("Contract errored: " + statusJson);
|
||||
}
|
||||
|
||||
if (DateTime.UtcNow - waitStart > timeout)
|
||||
{
|
||||
FrameworkAssert.Fail($"Contract did not reach '{desiredState}' within {Time.FormatDuration(timeout)} timeout. {statusJson}");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private void LogSubmittedDuration()
|
||||
{
|
||||
Log($"Pending to Submitted in {Time.FormatDuration(PendingToSubmitted)} " +
|
||||
$"( < {Time.FormatDuration(gracePeriod)})");
|
||||
}
|
||||
|
||||
private void LogStartedDuration()
|
||||
{
|
||||
Log($"Submitted to Started in {Time.FormatDuration(SubmittedToStarted)} " +
|
||||
$"( < {Time.FormatDuration(Purchase.Expiry + gracePeriod)})");
|
||||
}
|
||||
|
||||
private void LogFinishedDuration()
|
||||
{
|
||||
Log($"Submitted to Finished in {Time.FormatDuration(SubmittedToFinished)} " +
|
||||
$"( < {Time.FormatDuration(Purchase.Duration + gracePeriod)})");
|
||||
}
|
||||
|
||||
private void AssertDuration(TimeSpan? span, TimeSpan max, string message)
|
||||
{
|
||||
if (span == null) throw new ArgumentNullException(nameof(MarketplaceAccess) + ": " + message + " (IsNull)");
|
||||
if (span.Value.TotalDays >= max.TotalSeconds)
|
||||
{
|
||||
throw new Exception(nameof(MarketplaceAccess) +
|
||||
$": Duration out of range. Max: {Time.FormatDuration(max)} but was: {Time.FormatDuration(span.Value)} " +
|
||||
message);
|
||||
}
|
||||
}
|
||||
|
||||
private void Log(string msg)
|
||||
{
|
||||
log.Log($"[{PurchaseId}] {msg}");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,128 @@
|
||||
using CodexPlugin.OverwatchSupport.LineConverters;
|
||||
using KubernetesWorkflow;
|
||||
using OverwatchTranscript;
|
||||
using Utils;
|
||||
|
||||
namespace CodexPlugin.OverwatchSupport
|
||||
{
|
||||
public class CodexLogConverter
|
||||
{
|
||||
private readonly ITranscriptWriter writer;
|
||||
private readonly CodexTranscriptWriterConfig config;
|
||||
private readonly IdentityMap identityMap;
|
||||
|
||||
public CodexLogConverter(ITranscriptWriter writer, CodexTranscriptWriterConfig config, IdentityMap identityMap)
|
||||
{
|
||||
this.writer = writer;
|
||||
this.config = config;
|
||||
this.identityMap = identityMap;
|
||||
}
|
||||
|
||||
public void ProcessLog(IDownloadedLog log)
|
||||
{
|
||||
var name = DetermineName(log);
|
||||
var identityIndex = identityMap.GetIndex(name);
|
||||
var runner = new ConversionRunner(writer, config, identityMap, identityIndex);
|
||||
runner.Run(log);
|
||||
}
|
||||
|
||||
private string DetermineName(IDownloadedLog log)
|
||||
{
|
||||
// Expected string:
|
||||
// Downloading container log for '<Downloader1>'
|
||||
var nameLine = log.FindLinesThatContain("Downloading container log for").First();
|
||||
return Str.Between(nameLine, "'", "'");
|
||||
}
|
||||
}
|
||||
|
||||
public class ConversionRunner
|
||||
{
|
||||
private readonly ITranscriptWriter writer;
|
||||
private readonly IdentityMap nameIdMap;
|
||||
private readonly int nodeIdentityIndex;
|
||||
private readonly ILineConverter[] converters;
|
||||
|
||||
public ConversionRunner(ITranscriptWriter writer, CodexTranscriptWriterConfig config, IdentityMap nameIdMap, int nodeIdentityIndex)
|
||||
{
|
||||
this.nodeIdentityIndex = nodeIdentityIndex;
|
||||
this.writer = writer;
|
||||
this.nameIdMap = nameIdMap;
|
||||
|
||||
converters = CreateConverters(config).ToArray();
|
||||
}
|
||||
|
||||
private IEnumerable<ILineConverter> CreateConverters(CodexTranscriptWriterConfig config)
|
||||
{
|
||||
if (config.IncludeBlockReceivedEvents)
|
||||
{
|
||||
yield return new BlockReceivedLineConverter();
|
||||
}
|
||||
yield return new BootstrapLineConverter();
|
||||
yield return new DialSuccessfulLineConverter();
|
||||
yield return new PeerDroppedLineConverter();
|
||||
}
|
||||
|
||||
public void Run(IDownloadedLog log)
|
||||
{
|
||||
log.IterateLines(line =>
|
||||
{
|
||||
foreach (var converter in converters)
|
||||
{
|
||||
ProcessLine(line, converter);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
private void AddEvent(DateTime utc, Action<OverwatchCodexEvent> action)
|
||||
{
|
||||
var e = new OverwatchCodexEvent
|
||||
{
|
||||
NodeIdentity = nodeIdentityIndex,
|
||||
};
|
||||
action(e);
|
||||
|
||||
e.Write(utc, writer);
|
||||
}
|
||||
|
||||
private void ProcessLine(string line, ILineConverter converter)
|
||||
{
|
||||
if (!line.Contains(converter.Interest)) return;
|
||||
|
||||
var codexLine = CodexLogLine.Parse(line);
|
||||
|
||||
if (codexLine == null) throw new Exception("Unable to parse required line");
|
||||
EnsureFullIds(codexLine);
|
||||
|
||||
converter.Process(codexLine, (action) =>
|
||||
{
|
||||
AddEvent(codexLine.TimestampUtc, action);
|
||||
});
|
||||
}
|
||||
|
||||
private void EnsureFullIds(CodexLogLine codexLine)
|
||||
{
|
||||
// The issue is: node IDs occure both in full and short version.
|
||||
// Downstream tools will assume that a node ID string-equals its own ID.
|
||||
// So we replace all shortened IDs we can find with their full ones.
|
||||
|
||||
// Usually, the shortID appears as the entire string of an attribute:
|
||||
// "peerId=123*567890"
|
||||
// But sometimes, it is part of a larger string:
|
||||
// "thing=abc:123*567890,def"
|
||||
|
||||
foreach (var pair in codexLine.Attributes)
|
||||
{
|
||||
if (pair.Value.Contains("*"))
|
||||
{
|
||||
codexLine.Attributes[pair.Key] = nameIdMap.ReplaceShortIds(pair.Value);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
public interface ILineConverter
|
||||
{
|
||||
string Interest { get; }
|
||||
void Process(CodexLogLine line, Action<Action<OverwatchCodexEvent>> addEvent);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,191 @@
|
||||
using CodexPlugin.Hooks;
|
||||
using GethPlugin;
|
||||
using OverwatchTranscript;
|
||||
using Utils;
|
||||
|
||||
namespace CodexPlugin.OverwatchSupport
|
||||
{
|
||||
public class CodexNodeTranscriptWriter : ICodexNodeHooks
|
||||
{
|
||||
private readonly ITranscriptWriter writer;
|
||||
private readonly IdentityMap identityMap;
|
||||
private readonly string name;
|
||||
private int identityIndex = -1;
|
||||
private readonly List<(DateTime, OverwatchCodexEvent)> pendingEvents = new List<(DateTime, OverwatchCodexEvent)>();
|
||||
|
||||
public CodexNodeTranscriptWriter(ITranscriptWriter writer, IdentityMap identityMap, string name)
|
||||
{
|
||||
this.writer = writer;
|
||||
this.identityMap = identityMap;
|
||||
this.name = name;
|
||||
}
|
||||
|
||||
public void OnNodeStarting(DateTime startUtc, string image, EthAccount? ethAccount)
|
||||
{
|
||||
WriteCodexEvent(startUtc, e =>
|
||||
{
|
||||
e.NodeStarting = new NodeStartingEvent
|
||||
{
|
||||
Image = image,
|
||||
EthAddress = ethAccount != null ? ethAccount.ToString() : ""
|
||||
};
|
||||
});
|
||||
}
|
||||
|
||||
public void OnNodeStarted(string peerId, string nodeId)
|
||||
{
|
||||
if (string.IsNullOrEmpty(peerId) || string.IsNullOrEmpty(nodeId))
|
||||
{
|
||||
throw new Exception("Node started - peerId and/or nodeId unknown.");
|
||||
}
|
||||
|
||||
identityMap.Add(name, peerId, nodeId);
|
||||
identityIndex = identityMap.GetIndex(name);
|
||||
|
||||
WriteCodexEvent(e =>
|
||||
{
|
||||
e.NodeStarted = new NodeStartedEvent
|
||||
{
|
||||
};
|
||||
});
|
||||
}
|
||||
|
||||
public void OnNodeStopping()
|
||||
{
|
||||
WriteCodexEvent(e =>
|
||||
{
|
||||
e.NodeStopping = new NodeStoppingEvent
|
||||
{
|
||||
};
|
||||
});
|
||||
}
|
||||
|
||||
public void OnFileDownloading(ContentId cid)
|
||||
{
|
||||
WriteCodexEvent(e =>
|
||||
{
|
||||
e.FileDownloading = new FileDownloadingEvent
|
||||
{
|
||||
Cid = cid.Id
|
||||
};
|
||||
});
|
||||
}
|
||||
|
||||
public void OnFileDownloaded(ByteSize size, ContentId cid)
|
||||
{
|
||||
WriteCodexEvent(e =>
|
||||
{
|
||||
e.FileDownloaded = new FileDownloadedEvent
|
||||
{
|
||||
Cid = cid.Id,
|
||||
ByteSize = size.SizeInBytes
|
||||
};
|
||||
});
|
||||
}
|
||||
|
||||
public void OnFileUploading(string uid, ByteSize size)
|
||||
{
|
||||
WriteCodexEvent(e =>
|
||||
{
|
||||
e.FileUploading = new FileUploadingEvent
|
||||
{
|
||||
UniqueId = uid,
|
||||
ByteSize = size.SizeInBytes
|
||||
};
|
||||
});
|
||||
}
|
||||
|
||||
public void OnFileUploaded(string uid, ByteSize size, ContentId cid)
|
||||
{
|
||||
WriteCodexEvent(e =>
|
||||
{
|
||||
e.FileUploaded = new FileUploadedEvent
|
||||
{
|
||||
UniqueId = uid,
|
||||
Cid = cid.Id,
|
||||
ByteSize = size.SizeInBytes
|
||||
};
|
||||
});
|
||||
}
|
||||
|
||||
public void OnStorageContractSubmitted(StoragePurchaseContract storagePurchaseContract)
|
||||
{
|
||||
WriteCodexEvent(e =>
|
||||
{
|
||||
e.StorageContractSubmitted = new StorageContractSubmittedEvent
|
||||
{
|
||||
PurchaseId = storagePurchaseContract.PurchaseId,
|
||||
PurchaseRequest = storagePurchaseContract.Purchase
|
||||
};
|
||||
});
|
||||
}
|
||||
|
||||
public void OnStorageContractUpdated(StoragePurchase purchaseStatus)
|
||||
{
|
||||
WriteCodexEvent(e =>
|
||||
{
|
||||
e.StorageContractUpdated = new StorageContractUpdatedEvent
|
||||
{
|
||||
StoragePurchase = purchaseStatus
|
||||
};
|
||||
});
|
||||
}
|
||||
|
||||
public void OnStorageAvailabilityCreated(StorageAvailability response)
|
||||
{
|
||||
WriteCodexEvent(e =>
|
||||
{
|
||||
e.StorageAvailabilityCreated = new StorageAvailabilityCreatedEvent
|
||||
{
|
||||
StorageAvailability = response
|
||||
};
|
||||
});
|
||||
}
|
||||
|
||||
private void WriteCodexEvent(Action<OverwatchCodexEvent> action)
|
||||
{
|
||||
WriteCodexEvent(DateTime.UtcNow, action);
|
||||
}
|
||||
|
||||
private void WriteCodexEvent(DateTime utc, Action<OverwatchCodexEvent> action)
|
||||
{
|
||||
var e = new OverwatchCodexEvent
|
||||
{
|
||||
NodeIdentity = identityIndex
|
||||
};
|
||||
|
||||
action(e);
|
||||
|
||||
if (identityIndex < 0)
|
||||
{
|
||||
// If we don't know our id, don't write the events yet.
|
||||
AddToCache(utc, e);
|
||||
}
|
||||
else
|
||||
{
|
||||
e.Write(utc, writer);
|
||||
|
||||
// Write any events that we cached when we didn't have our id yet.
|
||||
WriteAndClearCache();
|
||||
}
|
||||
}
|
||||
|
||||
private void AddToCache(DateTime utc, OverwatchCodexEvent e)
|
||||
{
|
||||
pendingEvents.Add((utc, e));
|
||||
}
|
||||
|
||||
private void WriteAndClearCache()
|
||||
{
|
||||
if (pendingEvents.Any())
|
||||
{
|
||||
foreach (var pair in pendingEvents)
|
||||
{
|
||||
pair.Item2.NodeIdentity = identityIndex;
|
||||
pair.Item2.Write(pair.Item1, writer);
|
||||
}
|
||||
pendingEvents.Clear();
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,91 @@
|
||||
using CodexPlugin.Hooks;
|
||||
using KubernetesWorkflow;
|
||||
using Logging;
|
||||
using OverwatchTranscript;
|
||||
using Utils;
|
||||
|
||||
namespace CodexPlugin.OverwatchSupport
|
||||
{
|
||||
public class CodexTranscriptWriter : ICodexHooksProvider
|
||||
{
|
||||
private const string CodexHeaderKey = "cdx_h";
|
||||
private readonly ILog log;
|
||||
private readonly CodexTranscriptWriterConfig config;
|
||||
private readonly ITranscriptWriter writer;
|
||||
private readonly CodexLogConverter converter;
|
||||
private readonly IdentityMap identityMap = new IdentityMap();
|
||||
private readonly KademliaPositionFinder positionFinder = new KademliaPositionFinder();
|
||||
|
||||
public CodexTranscriptWriter(ILog log, CodexTranscriptWriterConfig config, ITranscriptWriter transcriptWriter)
|
||||
{
|
||||
this.log = log;
|
||||
this.config = config;
|
||||
writer = transcriptWriter;
|
||||
converter = new CodexLogConverter(writer, config, identityMap);
|
||||
}
|
||||
|
||||
public void Finalize(string outputFilepath)
|
||||
{
|
||||
log.Log("Finalizing Codex transcript...");
|
||||
|
||||
writer.AddHeader(CodexHeaderKey, CreateCodexHeader());
|
||||
writer.Write(outputFilepath);
|
||||
|
||||
log.Log("Done");
|
||||
}
|
||||
|
||||
public ICodexNodeHooks CreateHooks(string nodeName)
|
||||
{
|
||||
nodeName = Str.Between(nodeName, "'", "'");
|
||||
return new CodexNodeTranscriptWriter(writer, identityMap, nodeName);
|
||||
}
|
||||
|
||||
public void IncludeFile(string filepath)
|
||||
{
|
||||
writer.IncludeArtifact(filepath);
|
||||
}
|
||||
|
||||
public void ProcessLogs(IDownloadedLog[] downloadedLogs)
|
||||
{
|
||||
foreach (var l in downloadedLogs)
|
||||
{
|
||||
log.Log("Include artifact: " + l.GetFilepath());
|
||||
writer.IncludeArtifact(l.GetFilepath());
|
||||
|
||||
// Not all of these logs are necessarily Codex logs.
|
||||
// Check, and process only the Codex ones.
|
||||
if (IsCodexLog(l))
|
||||
{
|
||||
log.Log("Processing Codex log: " + l.GetFilepath());
|
||||
converter.ProcessLog(l);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
public void AddResult(bool success, string result)
|
||||
{
|
||||
writer.Add(DateTime.UtcNow, new OverwatchCodexEvent
|
||||
{
|
||||
NodeIdentity = -1,
|
||||
ScenarioFinished = new ScenarioFinishedEvent
|
||||
{
|
||||
Success = success,
|
||||
Result = result
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
private OverwatchCodexHeader CreateCodexHeader()
|
||||
{
|
||||
return new OverwatchCodexHeader
|
||||
{
|
||||
Nodes = positionFinder.DeterminePositions(identityMap.Get())
|
||||
};
|
||||
}
|
||||
|
||||
private bool IsCodexLog(IDownloadedLog log)
|
||||
{
|
||||
return log.GetLinesContaining("Run Codex node").Any();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,12 @@
|
||||
namespace CodexPlugin.OverwatchSupport
|
||||
{
|
||||
public class CodexTranscriptWriterConfig
|
||||
{
|
||||
public CodexTranscriptWriterConfig(bool includeBlockReceivedEvents)
|
||||
{
|
||||
IncludeBlockReceivedEvents = includeBlockReceivedEvents;
|
||||
}
|
||||
|
||||
public bool IncludeBlockReceivedEvents { get; }
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,63 @@
|
||||
namespace CodexPlugin.OverwatchSupport
|
||||
{
|
||||
public class IdentityMap
|
||||
{
|
||||
private readonly List<CodexNodeIdentity> nodes = new List<CodexNodeIdentity>();
|
||||
private readonly Dictionary<string, int> nameIndexMap = new Dictionary<string, int>();
|
||||
private readonly Dictionary<string, string> shortToLong = new Dictionary<string, string>();
|
||||
|
||||
public void Add(string name, string peerId, string nodeId)
|
||||
{
|
||||
Add(new CodexNodeIdentity
|
||||
{
|
||||
Name = name,
|
||||
PeerId = peerId,
|
||||
NodeId = nodeId
|
||||
});
|
||||
|
||||
nameIndexMap.Add(name, nameIndexMap.Count);
|
||||
}
|
||||
|
||||
public void Add(CodexNodeIdentity identity)
|
||||
{
|
||||
if (string.IsNullOrWhiteSpace(identity.Name)) throw new Exception("Name required");
|
||||
if (string.IsNullOrWhiteSpace(identity.PeerId) || identity.PeerId.Length < 11) throw new Exception("PeerId invalid");
|
||||
if (string.IsNullOrWhiteSpace(identity.NodeId) || identity.NodeId.Length < 11) throw new Exception("NodeId invalid");
|
||||
|
||||
nodes.Add(identity);
|
||||
|
||||
shortToLong.Add(CodexUtils.ToShortId(identity.PeerId), identity.PeerId);
|
||||
shortToLong.Add(CodexUtils.ToNodeIdShortId(identity.NodeId), identity.NodeId);
|
||||
}
|
||||
|
||||
public CodexNodeIdentity[] Get()
|
||||
{
|
||||
return nodes.ToArray();
|
||||
}
|
||||
|
||||
public int GetIndex(string name)
|
||||
{
|
||||
return nameIndexMap[name];
|
||||
}
|
||||
|
||||
public CodexNodeIdentity GetId(string name)
|
||||
{
|
||||
return nodes.Single(n => n.Name == name);
|
||||
}
|
||||
|
||||
public string ReplaceShortIds(string value)
|
||||
{
|
||||
var result = value;
|
||||
foreach (var pair in shortToLong)
|
||||
{
|
||||
result = result.Replace(pair.Key, pair.Value);
|
||||
}
|
||||
return result;
|
||||
}
|
||||
|
||||
public int Size
|
||||
{
|
||||
get { return nodes.Count; }
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,68 @@
|
||||
using System.Numerics;
|
||||
using Utils;
|
||||
using YamlDotNet.Core.Tokens;
|
||||
|
||||
namespace CodexPlugin.OverwatchSupport
|
||||
{
|
||||
public class KademliaPositionFinder
|
||||
{
|
||||
public CodexNodeIdentity[] DeterminePositions(CodexNodeIdentity[] identities)
|
||||
{
|
||||
var zero = identities.First();
|
||||
var distances = CalculateDistances(zero, identities);
|
||||
|
||||
var maxDistance = distances.Values.Max();
|
||||
CalculateNormalizedPositions(distances, maxDistance);
|
||||
|
||||
return identities;
|
||||
}
|
||||
|
||||
private Dictionary<CodexNodeIdentity, BigInteger> CalculateDistances(CodexNodeIdentity zero, CodexNodeIdentity[] identities)
|
||||
{
|
||||
var result = new Dictionary<CodexNodeIdentity, BigInteger>();
|
||||
foreach (var id in identities.Skip(1))
|
||||
{
|
||||
result.Add(id, GetDistance(zero.NodeId, id.NodeId));
|
||||
}
|
||||
return result;
|
||||
}
|
||||
|
||||
private BigInteger GetDistance(string id1, string id2)
|
||||
{
|
||||
var one = BigInteger.Parse(id1, System.Globalization.NumberStyles.HexNumber).ToByteArray();
|
||||
var two = BigInteger.Parse(id2, System.Globalization.NumberStyles.HexNumber).ToByteArray();
|
||||
|
||||
var x = Xor(one, two);
|
||||
return new BigInteger(x, isUnsigned: true);
|
||||
}
|
||||
|
||||
private byte[] Xor(byte[] one, byte[] two)
|
||||
{
|
||||
if (one.Length != two.Length) throw new Exception("Not equal length");
|
||||
|
||||
var result = new byte[one.Length];
|
||||
for (int i = 0; i < one.Length; i++)
|
||||
{
|
||||
uint a = one[i];
|
||||
uint b = two[i];
|
||||
uint c = (a ^ b);
|
||||
result[i] = (byte)c;
|
||||
}
|
||||
return result;
|
||||
}
|
||||
|
||||
private void CalculateNormalizedPositions(Dictionary<CodexNodeIdentity, BigInteger> distances, BigInteger maxDistance)
|
||||
{
|
||||
foreach (var pair in distances)
|
||||
{
|
||||
pair.Key.KademliaNormalizedPosition = DeterminePosition(pair.Value, maxDistance);
|
||||
}
|
||||
}
|
||||
|
||||
private float DeterminePosition(BigInteger value, BigInteger maxDistance)
|
||||
{
|
||||
var f = (value * 10000) / maxDistance;
|
||||
return ((float)f) / 10000.0f;
|
||||
}
|
||||
}
|
||||
}
|
||||
+44
@@ -0,0 +1,44 @@
|
||||
namespace CodexPlugin.OverwatchSupport.LineConverters
|
||||
{
|
||||
public class BlockReceivedLineConverter : ILineConverter
|
||||
{
|
||||
public string Interest => "Received blocks from peer";
|
||||
|
||||
public void Process(CodexLogLine line, Action<Action<OverwatchCodexEvent>> addEvent)
|
||||
{
|
||||
var peer = line.Attributes["peer"];
|
||||
var blockAddresses = line.Attributes["blocks"];
|
||||
|
||||
SplitBlockAddresses(blockAddresses, address =>
|
||||
{
|
||||
addEvent(e =>
|
||||
{
|
||||
e.BlockReceived = new BlockReceivedEvent
|
||||
{
|
||||
SenderPeerId = peer,
|
||||
BlockAddress = address
|
||||
};
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
private void SplitBlockAddresses(string blockAddresses, Action<string> onBlockAddress)
|
||||
{
|
||||
// Single line can contain multiple block addresses.
|
||||
var tokens = blockAddresses.Split(",", StringSplitOptions.RemoveEmptyEntries).ToList();
|
||||
while (tokens.Count > 0)
|
||||
{
|
||||
if (tokens.Count == 1)
|
||||
{
|
||||
onBlockAddress(tokens[0]);
|
||||
return;
|
||||
}
|
||||
|
||||
var blockAddress = $"{tokens[0]}, {tokens[1]}";
|
||||
tokens.RemoveRange(0, 2);
|
||||
|
||||
onBlockAddress(blockAddress);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,71 @@
|
||||
namespace CodexPlugin.OverwatchSupport.LineConverters
|
||||
{
|
||||
public class BootstrapLineConverter : ILineConverter
|
||||
{
|
||||
private const string peerIdTag = "peerId: ";
|
||||
|
||||
public string Interest => "Starting codex node";
|
||||
|
||||
public void Process(CodexLogLine line, Action<Action<OverwatchCodexEvent>> addEvent)
|
||||
{
|
||||
// "(
|
||||
// configFile: none(InputFile),
|
||||
// logLevel: \"TRACE;warn:discv5,providers,manager,cache;warn:libp2p,multistream,switch,transport,tcptransport,semaphore,asyncstreamwrapper,lpstream,mplex,mplexchannel,noise,bufferstream,mplexcoder,secure,chronosstream,connection,connmanager,websock,ws-session,dialer,muxedupgrade,upgrade,identify;warn:contracts,clock;warn:serde,json,serialization\",
|
||||
// logFormat: auto,
|
||||
// metricsEnabled: false,
|
||||
// metricsAddress: 127.0.0.1,
|
||||
// metricsPort: 8008,
|
||||
// dataDir: datadir5,
|
||||
// circuitDir: /root/.cache/codex/circuits,
|
||||
// listenAddrs: @[/ip4/0.0.0.0/tcp/8081],
|
||||
// nat: 10.1.0.214,
|
||||
// discoveryIp: 0.0.0.0,
|
||||
// discoveryPort: 8080,
|
||||
// netPrivKeyFile: \"key\",
|
||||
// bootstrapNodes:
|
||||
// @[(envelope: (publicKey: secp256k1 key (0414380858330307a4a59e8a1643512f80680dedb8c541674f40382be71c26556cd65b7e1775ec9fabfaf6d58d562b88a3c8afb969f8cc256db20e4c4c9e1f70a6),
|
||||
// domain: \"libp2p-peer-record\",
|
||||
// payloadType: @[3, 1],
|
||||
// payload: @[10, 39, 0, 37, 8, 2, 18, 33, 2, 20, 56, 8, 88, 51, 3, 7, 164, 165, 158, 138, 22, 67, 81, 47, 128, 104, 13, 237, 184, 197, 65, 103, 79, 64, 56, 43, 231, 28, 38, 85, 108, 16, 254, 211, 141, 181, 6, 26, 11, 10, 9, 4, 10, 1, 0, 210, 145, 2, 31, 144],
|
||||
// signature: 3045022100FA846871D96EDCA579990244B1C590E16B57A7BDB817908A2B580043F8A0B0280220342825DBE577E83C22CD59AA9099DE8F8DC86A38C659F60E8B9255126CF37FDE),
|
||||
// data:
|
||||
// (peerId: 16Uiu2HAkvnbgwdB2NmNe1uWGJxE3Ep3sDHxNW95W2rTPs2vfxvou,
|
||||
// seqNo: 1721985534,
|
||||
// addresses: @[(address: /ip4/10.1.0.210/udp/8080)]
|
||||
// )
|
||||
// )],
|
||||
// maxPeers: 160,
|
||||
// agentString: \"Codex\",
|
||||
// apiBindAddress: \"0.0.0.0\",
|
||||
// apiPort: 30035,
|
||||
// apiCorsAllowedOrigin: none(string),
|
||||
// repoKind: fs,
|
||||
// storageQuota: 8589934592\'NByte,
|
||||
// blockTtl: 1d,
|
||||
// blockMaintenanceInterval: 10m,
|
||||
// blockMaintenanceNumberOfBlocks: 1000,
|
||||
// cacheSize: 0\'NByte,
|
||||
// logFile: none(string),
|
||||
// cmd: noCmd
|
||||
// )"
|
||||
|
||||
var config = line.Attributes["config"];
|
||||
|
||||
while (config.Contains(peerIdTag))
|
||||
{
|
||||
var openIndex = config.IndexOf(peerIdTag) + peerIdTag.Length;
|
||||
var closeIndex = config.IndexOf(",", openIndex);
|
||||
var bootPeerId = config.Substring(openIndex, closeIndex - openIndex);
|
||||
config = config.Substring(closeIndex);
|
||||
|
||||
addEvent(e =>
|
||||
{
|
||||
e.BootstrapConfig = new BootstrapConfigEvent
|
||||
{
|
||||
BootstrapPeerId = bootPeerId
|
||||
};
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
+20
@@ -0,0 +1,20 @@
|
||||
namespace CodexPlugin.OverwatchSupport.LineConverters
|
||||
{
|
||||
public class DialSuccessfulLineConverter : ILineConverter
|
||||
{
|
||||
public string Interest => "Dial successful";
|
||||
|
||||
public void Process(CodexLogLine line, Action<Action<OverwatchCodexEvent>> addEvent)
|
||||
{
|
||||
var peerId = line.Attributes["peerId"];
|
||||
|
||||
addEvent(e =>
|
||||
{
|
||||
e.DialSuccessful = new PeerDialSuccessfulEvent
|
||||
{
|
||||
TargetPeerId = peerId
|
||||
};
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
+20
@@ -0,0 +1,20 @@
|
||||
namespace CodexPlugin.OverwatchSupport.LineConverters
|
||||
{
|
||||
public class PeerDroppedLineConverter : ILineConverter
|
||||
{
|
||||
public string Interest => "Dropping peer";
|
||||
|
||||
public void Process(CodexLogLine line, Action<Action<OverwatchCodexEvent>> addEvent)
|
||||
{
|
||||
var peerId = line.Attributes["peer"];
|
||||
|
||||
addEvent(e =>
|
||||
{
|
||||
e.PeerDropped = new PeerDroppedEvent
|
||||
{
|
||||
DroppedPeerId = peerId
|
||||
};
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,164 @@
|
||||
using OverwatchTranscript;
|
||||
|
||||
namespace CodexPlugin.OverwatchSupport
|
||||
{
|
||||
[Serializable]
|
||||
public class OverwatchCodexHeader
|
||||
{
|
||||
public CodexNodeIdentity[] Nodes { get; set; } = Array.Empty<CodexNodeIdentity>();
|
||||
}
|
||||
|
||||
[Serializable]
|
||||
public class OverwatchCodexEvent
|
||||
{
|
||||
public int NodeIdentity { get; set; } = -1;
|
||||
public ScenarioFinishedEvent? ScenarioFinished { get; set; }
|
||||
public NodeStartingEvent? NodeStarting { get; set; }
|
||||
public NodeStartedEvent? NodeStarted { get; set; }
|
||||
public NodeStoppingEvent? NodeStopping { get; set; }
|
||||
public BootstrapConfigEvent? BootstrapConfig { get; set; }
|
||||
public FileUploadingEvent? FileUploading { get; set; }
|
||||
public FileUploadedEvent? FileUploaded { get; set; }
|
||||
public FileDownloadingEvent? FileDownloading { get; set; }
|
||||
public FileDownloadedEvent? FileDownloaded { get; set; }
|
||||
public BlockReceivedEvent? BlockReceived { get; set; }
|
||||
public PeerDialSuccessfulEvent? DialSuccessful { get; set; }
|
||||
public PeerDroppedEvent? PeerDropped { get; set; }
|
||||
public StorageContractSubmittedEvent? StorageContractSubmitted { get; set; }
|
||||
public StorageContractUpdatedEvent? StorageContractUpdated { get; set; }
|
||||
public StorageAvailabilityCreatedEvent? StorageAvailabilityCreated { get; set; }
|
||||
|
||||
public void Write(DateTime utc, ITranscriptWriter writer)
|
||||
{
|
||||
if (NodeIdentity == -1 && ScenarioFinished == null)
|
||||
{
|
||||
throw new Exception("NodeIdentity not set, and event is not ScenarioFinished.");
|
||||
}
|
||||
if (AllNull()) throw new Exception("No event data was set");
|
||||
|
||||
writer.Add(utc, this);
|
||||
}
|
||||
|
||||
private bool AllNull()
|
||||
{
|
||||
var props = GetType()
|
||||
.GetProperties(System.Reflection.BindingFlags.Public | System.Reflection.BindingFlags.Instance)
|
||||
.Where(p => p.PropertyType != typeof(string)).ToArray();
|
||||
|
||||
return props.All(p => p.GetValue(this) == null);
|
||||
}
|
||||
}
|
||||
|
||||
[Serializable]
|
||||
public class CodexNodeIdentity
|
||||
{
|
||||
public string Name { get; set; } = string.Empty;
|
||||
public string PeerId { get; set; } = string.Empty;
|
||||
public string NodeId { get; set; } = string.Empty;
|
||||
public float KademliaNormalizedPosition { get; set; } = 0.0f;
|
||||
}
|
||||
|
||||
#region Scenario Generated Events
|
||||
|
||||
[Serializable]
|
||||
public class ScenarioFinishedEvent
|
||||
{
|
||||
public bool Success { get; set; }
|
||||
public string Result { get; set; } = string.Empty;
|
||||
}
|
||||
|
||||
[Serializable]
|
||||
public class NodeStartingEvent
|
||||
{
|
||||
public string Image { get; set; } = string.Empty;
|
||||
public string EthAddress { get; set; } = string.Empty;
|
||||
}
|
||||
|
||||
[Serializable]
|
||||
public class NodeStartedEvent
|
||||
{
|
||||
}
|
||||
|
||||
[Serializable]
|
||||
public class NodeStoppingEvent
|
||||
{
|
||||
}
|
||||
|
||||
[Serializable]
|
||||
public class BootstrapConfigEvent
|
||||
{
|
||||
public string BootstrapPeerId { get; set; } = string.Empty;
|
||||
}
|
||||
|
||||
[Serializable]
|
||||
public class FileUploadingEvent
|
||||
{
|
||||
public string UniqueId { get; set;} = string.Empty;
|
||||
public long ByteSize { get; set; }
|
||||
}
|
||||
|
||||
[Serializable]
|
||||
public class FileDownloadingEvent
|
||||
{
|
||||
public string Cid { get; set; } = string.Empty;
|
||||
}
|
||||
|
||||
[Serializable]
|
||||
public class FileUploadedEvent
|
||||
{
|
||||
public string UniqueId { get; set; } = string.Empty;
|
||||
public long ByteSize { get; set; }
|
||||
public string Cid { get; set; } = string.Empty;
|
||||
}
|
||||
|
||||
[Serializable]
|
||||
public class FileDownloadedEvent
|
||||
{
|
||||
public string Cid { get; set; } = string.Empty;
|
||||
public long ByteSize { get; set; }
|
||||
}
|
||||
|
||||
[Serializable]
|
||||
public class StorageAvailabilityCreatedEvent
|
||||
{
|
||||
public StorageAvailability StorageAvailability { get; set; } = null!;
|
||||
}
|
||||
|
||||
[Serializable]
|
||||
public class StorageContractUpdatedEvent
|
||||
{
|
||||
public StoragePurchase StoragePurchase { get; set; } = null!;
|
||||
}
|
||||
|
||||
[Serializable]
|
||||
public class StorageContractSubmittedEvent
|
||||
{
|
||||
public string PurchaseId { get; set; } = string.Empty;
|
||||
public StoragePurchaseRequest PurchaseRequest { get; set; } = null!;
|
||||
}
|
||||
|
||||
#endregion
|
||||
|
||||
#region Codex Generated Events
|
||||
|
||||
[Serializable]
|
||||
public class BlockReceivedEvent
|
||||
{
|
||||
public string BlockAddress { get; set; } = string.Empty;
|
||||
public string SenderPeerId { get; set; } = string.Empty;
|
||||
}
|
||||
|
||||
[Serializable]
|
||||
public class PeerDialSuccessfulEvent
|
||||
{
|
||||
public string TargetPeerId { get; set; } = string.Empty;
|
||||
}
|
||||
|
||||
[Serializable]
|
||||
public class PeerDroppedEvent
|
||||
{
|
||||
public string DroppedPeerId { get; set; } = string.Empty;
|
||||
}
|
||||
|
||||
#endregion
|
||||
}
|
||||
@@ -0,0 +1,149 @@
|
||||
using CodexPlugin.Hooks;
|
||||
using Logging;
|
||||
using Newtonsoft.Json;
|
||||
using Utils;
|
||||
|
||||
namespace CodexPlugin
|
||||
{
|
||||
public interface IStoragePurchaseContract
|
||||
{
|
||||
string PurchaseId { get; }
|
||||
StoragePurchaseRequest Purchase { get; }
|
||||
ContentId ContentId { get; }
|
||||
void WaitForStorageContractSubmitted();
|
||||
void WaitForStorageContractStarted();
|
||||
void WaitForStorageContractFinished();
|
||||
}
|
||||
|
||||
public class StoragePurchaseContract : IStoragePurchaseContract
|
||||
{
|
||||
private readonly ILog log;
|
||||
private readonly CodexAccess codexAccess;
|
||||
private readonly ICodexNodeHooks hooks;
|
||||
private readonly TimeSpan gracePeriod = TimeSpan.FromSeconds(30);
|
||||
private readonly DateTime contractPendingUtc = DateTime.UtcNow;
|
||||
private DateTime? contractSubmittedUtc = DateTime.UtcNow;
|
||||
private DateTime? contractStartedUtc;
|
||||
private DateTime? contractFinishedUtc;
|
||||
|
||||
public StoragePurchaseContract(ILog log, CodexAccess codexAccess, string purchaseId, StoragePurchaseRequest purchase, ICodexNodeHooks hooks)
|
||||
{
|
||||
this.log = log;
|
||||
this.codexAccess = codexAccess;
|
||||
PurchaseId = purchaseId;
|
||||
Purchase = purchase;
|
||||
this.hooks = hooks;
|
||||
ContentId = new ContentId(codexAccess.GetPurchaseStatus(purchaseId).Request.Content.Cid);
|
||||
}
|
||||
|
||||
public string PurchaseId { get; }
|
||||
public StoragePurchaseRequest Purchase { get; }
|
||||
public ContentId ContentId { get; }
|
||||
|
||||
public TimeSpan? PendingToSubmitted => contractSubmittedUtc - contractPendingUtc;
|
||||
public TimeSpan? SubmittedToStarted => contractStartedUtc - contractSubmittedUtc;
|
||||
public TimeSpan? SubmittedToFinished => contractFinishedUtc - contractSubmittedUtc;
|
||||
|
||||
public void WaitForStorageContractSubmitted()
|
||||
{
|
||||
WaitForStorageContractState(gracePeriod, "submitted", sleep: 200);
|
||||
contractSubmittedUtc = DateTime.UtcNow;
|
||||
LogSubmittedDuration();
|
||||
AssertDuration(PendingToSubmitted, gracePeriod, nameof(PendingToSubmitted));
|
||||
}
|
||||
|
||||
public void WaitForStorageContractStarted()
|
||||
{
|
||||
var timeout = Purchase.Expiry + gracePeriod;
|
||||
|
||||
WaitForStorageContractState(timeout, "started");
|
||||
contractStartedUtc = DateTime.UtcNow;
|
||||
LogStartedDuration();
|
||||
AssertDuration(SubmittedToStarted, timeout, nameof(SubmittedToStarted));
|
||||
}
|
||||
|
||||
public void WaitForStorageContractFinished()
|
||||
{
|
||||
if (!contractStartedUtc.HasValue)
|
||||
{
|
||||
WaitForStorageContractStarted();
|
||||
}
|
||||
var currentContractTime = DateTime.UtcNow - contractSubmittedUtc!.Value;
|
||||
var timeout = (Purchase.Duration - currentContractTime) + gracePeriod;
|
||||
WaitForStorageContractState(timeout, "finished");
|
||||
contractFinishedUtc = DateTime.UtcNow;
|
||||
LogFinishedDuration();
|
||||
AssertDuration(SubmittedToFinished, timeout, nameof(SubmittedToFinished));
|
||||
}
|
||||
|
||||
public StoragePurchase GetPurchaseStatus(string purchaseId)
|
||||
{
|
||||
return codexAccess.GetPurchaseStatus(purchaseId);
|
||||
}
|
||||
|
||||
private void WaitForStorageContractState(TimeSpan timeout, string desiredState, int sleep = 1000)
|
||||
{
|
||||
var lastState = "";
|
||||
var waitStart = DateTime.UtcNow;
|
||||
|
||||
Log($"Waiting for {Time.FormatDuration(timeout)} to reach state '{desiredState}'.");
|
||||
while (lastState != desiredState)
|
||||
{
|
||||
Thread.Sleep(sleep);
|
||||
|
||||
var purchaseStatus = codexAccess.GetPurchaseStatus(PurchaseId);
|
||||
var statusJson = JsonConvert.SerializeObject(purchaseStatus);
|
||||
if (purchaseStatus != null && purchaseStatus.State != lastState)
|
||||
{
|
||||
lastState = purchaseStatus.State;
|
||||
log.Debug("Purchase status: " + statusJson);
|
||||
hooks.OnStorageContractUpdated(purchaseStatus);
|
||||
}
|
||||
|
||||
if (lastState == "errored")
|
||||
{
|
||||
FrameworkAssert.Fail("Contract errored: " + statusJson);
|
||||
}
|
||||
|
||||
if (DateTime.UtcNow - waitStart > timeout)
|
||||
{
|
||||
FrameworkAssert.Fail($"Contract did not reach '{desiredState}' within {Time.FormatDuration(timeout)} timeout. {statusJson}");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private void LogSubmittedDuration()
|
||||
{
|
||||
Log($"Pending to Submitted in {Time.FormatDuration(PendingToSubmitted)} " +
|
||||
$"( < {Time.FormatDuration(gracePeriod)})");
|
||||
}
|
||||
|
||||
private void LogStartedDuration()
|
||||
{
|
||||
Log($"Submitted to Started in {Time.FormatDuration(SubmittedToStarted)} " +
|
||||
$"( < {Time.FormatDuration(Purchase.Expiry + gracePeriod)})");
|
||||
}
|
||||
|
||||
private void LogFinishedDuration()
|
||||
{
|
||||
Log($"Submitted to Finished in {Time.FormatDuration(SubmittedToFinished)} " +
|
||||
$"( < {Time.FormatDuration(Purchase.Duration + gracePeriod)})");
|
||||
}
|
||||
|
||||
private void AssertDuration(TimeSpan? span, TimeSpan max, string message)
|
||||
{
|
||||
if (span == null) throw new ArgumentNullException(nameof(MarketplaceAccess) + ": " + message + " (IsNull)");
|
||||
if (span.Value.TotalDays >= max.TotalSeconds)
|
||||
{
|
||||
throw new Exception(nameof(MarketplaceAccess) +
|
||||
$": Duration out of range. Max: {Time.FormatDuration(max)} but was: {Time.FormatDuration(span.Value)} " +
|
||||
message);
|
||||
}
|
||||
}
|
||||
|
||||
private void Log(string msg)
|
||||
{
|
||||
log.Log($"[{PurchaseId}] {msg}");
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -3,24 +3,29 @@ using System.Text;
|
||||
|
||||
public static class Program
|
||||
{
|
||||
private const string OpenApiFile = "../CodexPlugin/openapi.yaml";
|
||||
private const string ClientFile = "../CodexPlugin/obj/openapiClient.cs";
|
||||
private const string Search = "<INSERT-OPENAPI-YAML-HASH>";
|
||||
private const string TargetFile = "ApiChecker.cs";
|
||||
private const string CodexPluginFolderName = "CodexPlugin";
|
||||
private const string ProjectPluginsFolderName = "ProjectPlugins";
|
||||
|
||||
public static void Main(string[] args)
|
||||
{
|
||||
Console.WriteLine("Injecting hash of 'openapi.yaml'...");
|
||||
|
||||
// Force client rebuild by deleting previous artifact.
|
||||
File.Delete(ClientFile);
|
||||
var root = FindCodexPluginFolder();
|
||||
Console.WriteLine("Located CodexPlugin: " + root);
|
||||
var openApiFile = Path.Combine(root, "openapi.yaml");
|
||||
var clientFile = Path.Combine(root, "obj", "openapiClient.cs");
|
||||
var targetFile = Path.Combine(root, "ApiChecker.cs");
|
||||
|
||||
var hash = CreateHash();
|
||||
// Force client rebuild by deleting previous artifact.
|
||||
File.Delete(clientFile);
|
||||
|
||||
var hash = CreateHash(openApiFile);
|
||||
// This hash is used to verify that the Codex docker image being used is compatible
|
||||
// with the openapi.yaml being used by the Codex plugin.
|
||||
// If the openapi.yaml files don't match, an exception is thrown.
|
||||
|
||||
SearchAndInject(hash);
|
||||
SearchAndInject(hash, targetFile);
|
||||
|
||||
// This program runs as the pre-build trigger for "CodexPlugin".
|
||||
// You might be wondering why this work isn't done by a shell script.
|
||||
@@ -33,9 +38,39 @@ public static class Program
|
||||
Console.WriteLine("Done!");
|
||||
}
|
||||
|
||||
private static string CreateHash()
|
||||
private static string FindCodexPluginFolder()
|
||||
{
|
||||
var file = File.ReadAllText(OpenApiFile);
|
||||
var current = Directory.GetCurrentDirectory();
|
||||
|
||||
while (true)
|
||||
{
|
||||
var localFolders = Directory.GetDirectories(current);
|
||||
var projectPluginsFolders = localFolders.Where(l => l.EndsWith(ProjectPluginsFolderName)).ToArray();
|
||||
if (projectPluginsFolders.Length == 1)
|
||||
{
|
||||
return Path.Combine(projectPluginsFolders.Single(), CodexPluginFolderName);
|
||||
}
|
||||
var codexPluginFolders = localFolders.Where(l => l.EndsWith(CodexPluginFolderName)).ToArray();
|
||||
if (codexPluginFolders.Length == 1)
|
||||
{
|
||||
return codexPluginFolders.Single();
|
||||
}
|
||||
|
||||
var parent = Directory.GetParent(current);
|
||||
if (parent == null)
|
||||
{
|
||||
var msg = $"Unable to locate '{CodexPluginFolderName}' folder. Travelled up from: '{Directory.GetCurrentDirectory()}'";
|
||||
Console.WriteLine(msg);
|
||||
throw new Exception(msg);
|
||||
}
|
||||
|
||||
current = parent.FullName;
|
||||
}
|
||||
}
|
||||
|
||||
private static string CreateHash(string openApiFile)
|
||||
{
|
||||
var file = File.ReadAllText(openApiFile);
|
||||
var fileBytes = Encoding.ASCII.GetBytes(file
|
||||
.Replace(Environment.NewLine, ""));
|
||||
|
||||
@@ -44,11 +79,11 @@ public static class Program
|
||||
return BitConverter.ToString(hash);
|
||||
}
|
||||
|
||||
private static void SearchAndInject(string hash)
|
||||
private static void SearchAndInject(string hash, string targetFile)
|
||||
{
|
||||
var lines = File.ReadAllLines(TargetFile);
|
||||
var lines = File.ReadAllLines(targetFile);
|
||||
Inject(lines, hash);
|
||||
File.WriteAllLines(TargetFile, lines);
|
||||
File.WriteAllLines(targetFile, lines);
|
||||
}
|
||||
|
||||
private static void Inject(string[] lines, string hash)
|
||||
|
||||
@@ -24,5 +24,10 @@ namespace GethPlugin
|
||||
|
||||
return new EthAccount(ethAddress, account.PrivateKey);
|
||||
}
|
||||
|
||||
public override string ToString()
|
||||
{
|
||||
return EthAddress.ToString();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -15,9 +15,10 @@ namespace MetricsPlugin
|
||||
{
|
||||
RunningContainer = runningContainer;
|
||||
log = tools.GetLog();
|
||||
var address = RunningContainer.GetAddress(log, PrometheusContainerRecipe.PortTag);
|
||||
endpoint = tools
|
||||
.CreateHttp()
|
||||
.CreateEndpoint(RunningContainer.GetAddress(log, PrometheusContainerRecipe.PortTag), "/api/v1/");
|
||||
.CreateHttp(address.ToString())
|
||||
.CreateEndpoint(address, "/api/v1/");
|
||||
}
|
||||
|
||||
public RunningContainer RunningContainer { get; }
|
||||
|
||||
@@ -71,7 +71,7 @@ namespace ContinuousTests
|
||||
var address = new Address($"http://{serviceName}.{k8sNamespace}.svc.cluster.local", 9200);
|
||||
var baseUrl = "";
|
||||
|
||||
var http = tools.CreateHttp(client =>
|
||||
var http = tools.CreateHttp(address.ToString(), client =>
|
||||
{
|
||||
client.DefaultRequestHeaders.Add("kbn-xsrf", "reporting");
|
||||
});
|
||||
|
||||
@@ -4,6 +4,7 @@ using Utils;
|
||||
using Core;
|
||||
using CodexPlugin;
|
||||
using KubernetesWorkflow.Types;
|
||||
using KubernetesWorkflow;
|
||||
|
||||
namespace ContinuousTests
|
||||
{
|
||||
|
||||
@@ -1,7 +1,6 @@
|
||||
using CodexPlugin;
|
||||
using FileUtils;
|
||||
using Logging;
|
||||
using Newtonsoft.Json;
|
||||
using NUnit.Framework;
|
||||
using Utils;
|
||||
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
using DistTestCore;
|
||||
using CodexPlugin;
|
||||
using DistTestCore;
|
||||
using NUnit.Framework;
|
||||
using Utils;
|
||||
|
||||
@@ -16,59 +17,56 @@ namespace CodexTests.ScalabilityTests
|
||||
[Values(100, 1000)] int fileSize
|
||||
)
|
||||
{
|
||||
var hosts = StartCodex(numberOfHosts, s => s.WithLogLevel(CodexPlugin.CodexLogLevel.Trace));
|
||||
var hosts = StartCodex(numberOfHosts, s => s.WithName("host").WithLogLevel(CodexLogLevel.Trace));
|
||||
var file = GenerateTestFile(fileSize.MB());
|
||||
var cid = hosts[0].UploadFile(file);
|
||||
var tailOfManifestCid = cid.Id.Substring(cid.Id.Length - 6);
|
||||
|
||||
var uploadTasks = hosts.Select(h =>
|
||||
{
|
||||
return Task.Run(() =>
|
||||
{
|
||||
return h.UploadFile(file);
|
||||
});
|
||||
}).ToArray();
|
||||
|
||||
Task.WaitAll(uploadTasks);
|
||||
var cid = new ContentId(uploadTasks.Select(t => t.Result.Id).Distinct().Single());
|
||||
|
||||
var uploadLog = Ci.DownloadLog(hosts[0]);
|
||||
var expectedNumberOfBlocks = RoundUp(fileSize.MB().SizeInBytes, 64.KB().SizeInBytes) + 1; // +1 for manifest block.
|
||||
var blockCids = uploadLog
|
||||
.FindLinesThatContain("Putting block into network store")
|
||||
.FindLinesThatContain("Block Stored")
|
||||
.Select(s =>
|
||||
{
|
||||
var start = s.IndexOf("cid=") + 4;
|
||||
var end = s.IndexOf(" count=");
|
||||
var len = end - start;
|
||||
return s.Substring(start, len);
|
||||
var line = CodexLogLine.Parse(s)!;
|
||||
return line.Attributes["cid"];
|
||||
})
|
||||
.ToArray();
|
||||
|
||||
Assert.That(blockCids.Length, Is.EqualTo(expectedNumberOfBlocks));
|
||||
|
||||
foreach (var h in hosts) h.DownloadContent(cid);
|
||||
|
||||
var client = StartCodex(s => s.WithLogLevel(CodexPlugin.CodexLogLevel.Trace));
|
||||
|
||||
var client = StartCodex(s => s.WithName("client").WithLogLevel(CodexLogLevel.Trace));
|
||||
var resultFile = client.DownloadContent(cid);
|
||||
resultFile!.AssertIsEqual(file);
|
||||
|
||||
var downloadLog = Ci.DownloadLog(client);
|
||||
var host = string.Empty;
|
||||
var blockCidHostMap = new Dictionary<string, string>();
|
||||
downloadLog.IterateLines(line =>
|
||||
var blockAddressHostMap = new Dictionary<string, List<string>>();
|
||||
downloadLog
|
||||
.IterateLines(s =>
|
||||
{
|
||||
if (line.Contains("peer=") && line.Contains(" len="))
|
||||
{
|
||||
var start = line.IndexOf("peer=") + 5;
|
||||
var end = line.IndexOf(" len=");
|
||||
var len = end - start;
|
||||
host = line.Substring(start, len);
|
||||
}
|
||||
else if (!string.IsNullOrEmpty(host) && line.Contains("Storing block with key"))
|
||||
{
|
||||
var start = line.IndexOf("cid=") + 4;
|
||||
var end = line.IndexOf(" count=");
|
||||
var len = end - start;
|
||||
var blockCid = line.Substring(start, len);
|
||||
var line = CodexLogLine.Parse(s)!;
|
||||
var peer = line.Attributes["peer"];
|
||||
var blockAddresses = line.Attributes["blocks"];
|
||||
|
||||
blockCidHostMap.Add(blockCid, host);
|
||||
host = string.Empty;
|
||||
}
|
||||
});
|
||||
AddBlockAddresses(peer, blockAddresses, blockAddressHostMap);
|
||||
|
||||
var totalFetched = blockCidHostMap.Count(p => !string.IsNullOrEmpty(p.Value));
|
||||
//PrintFullMap(blockCidHostMap);
|
||||
PrintOverview(blockCidHostMap);
|
||||
}, thatContain: "Received blocks from peer");
|
||||
|
||||
var totalFetched = blockAddressHostMap.Count(p => p.Value.Any());
|
||||
PrintFullMap(blockAddressHostMap);
|
||||
//PrintOverview(blockCidHostMap);
|
||||
|
||||
Log("Expected number of blocks: " + expectedNumberOfBlocks);
|
||||
Log("Total number of block CIDs found in dataset + manifest block: " + blockCids.Length);
|
||||
@@ -76,13 +74,47 @@ namespace CodexTests.ScalabilityTests
|
||||
Assert.That(totalFetched, Is.EqualTo(expectedNumberOfBlocks));
|
||||
}
|
||||
|
||||
private void PrintOverview(Dictionary<string, string> blockCidHostMap)
|
||||
private void AddBlockAddresses(string peer, string blockAddresses, Dictionary<string, List<string>> blockAddressHostMap)
|
||||
{
|
||||
// Single line can contain multiple block addresses.
|
||||
var tokens = blockAddresses.Split(",", StringSplitOptions.RemoveEmptyEntries).ToList();
|
||||
while (tokens.Count > 0)
|
||||
{
|
||||
if (tokens.Count == 1)
|
||||
{
|
||||
AddBlockAddress(peer, tokens[0], blockAddressHostMap);
|
||||
return;
|
||||
}
|
||||
|
||||
var blockAddress = $"{tokens[0]}, {tokens[1]}";
|
||||
tokens.RemoveRange(0, 2);
|
||||
|
||||
AddBlockAddress(peer, blockAddress, blockAddressHostMap);
|
||||
}
|
||||
}
|
||||
|
||||
private void AddBlockAddress(string peer, string blockAddress, Dictionary<string, List<string>> blockAddressHostMap)
|
||||
{
|
||||
if (blockAddressHostMap.ContainsKey(blockAddress))
|
||||
{
|
||||
blockAddressHostMap[blockAddress].Add(peer);
|
||||
}
|
||||
else
|
||||
{
|
||||
blockAddressHostMap[blockAddress] = new List<string> { peer };
|
||||
}
|
||||
}
|
||||
|
||||
private void PrintOverview(Dictionary<string, List<string>> blockAddressHostMap)
|
||||
{
|
||||
var overview = new Dictionary<string, int>();
|
||||
foreach (var pair in blockCidHostMap)
|
||||
foreach (var pair in blockAddressHostMap)
|
||||
{
|
||||
if (!overview.ContainsKey(pair.Value)) overview.Add(pair.Value, 1);
|
||||
else overview[pair.Value]++;
|
||||
foreach (var host in pair.Value)
|
||||
{
|
||||
if (!overview.ContainsKey(host)) overview.Add(host, 1);
|
||||
else overview[host]++;
|
||||
}
|
||||
}
|
||||
|
||||
Log("Blocks fetched per host:");
|
||||
@@ -92,19 +124,13 @@ namespace CodexTests.ScalabilityTests
|
||||
}
|
||||
}
|
||||
|
||||
private void PrintFullMap(Dictionary<string, string> blockCidHostMap)
|
||||
private void PrintFullMap(Dictionary<string, List<string>> blockAddressHostMap)
|
||||
{
|
||||
Log("Per block, host it was fetched from:");
|
||||
foreach (var pair in blockCidHostMap)
|
||||
foreach (var pair in blockAddressHostMap)
|
||||
{
|
||||
if (string.IsNullOrEmpty(pair.Value))
|
||||
{
|
||||
Log($"block: {pair.Key} = Not seen");
|
||||
}
|
||||
else
|
||||
{
|
||||
Log($"block: {pair.Key} = '{pair.Value}'");
|
||||
}
|
||||
var hostStr = $"[{string.Join(",", pair.Value)}]";
|
||||
Log($"blockAddress: {pair.Key} = {hostStr}");
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
using CodexContractsPlugin;
|
||||
using CodexContractsPlugin.Marketplace;
|
||||
using CodexPlugin;
|
||||
using FileUtils;
|
||||
using GethPlugin;
|
||||
using NUnit.Framework;
|
||||
using Utils;
|
||||
@@ -11,11 +12,21 @@ namespace CodexTests.BasicTests
|
||||
public class MarketplaceTests : AutoBootstrapDistTest
|
||||
{
|
||||
[Test]
|
||||
public void MarketplaceExample()
|
||||
[Combinatorial]
|
||||
[CreateTranscript(nameof(MarketplaceTests), includeBlockReceivedEvents: false)]
|
||||
public void MarketplaceExample(
|
||||
[Values(64)] int numBlocks,
|
||||
[Values(0)] int plusSizeKb,
|
||||
[Values(0)] int plusSizeBytes
|
||||
)
|
||||
{
|
||||
var hostInitialBalance = 234.TstWei();
|
||||
var clientInitialBalance = 100000.TstWei();
|
||||
var fileSize = 10.MB();
|
||||
var fileSize = new ByteSize(
|
||||
numBlocks * (64 * 1024) +
|
||||
plusSizeKb * 1024 +
|
||||
plusSizeBytes
|
||||
);
|
||||
|
||||
var geth = Ci.StartGethNode(s => s.IsMiner().WithName("disttest-geth"));
|
||||
var contracts = Ci.StartCodexContracts(geth);
|
||||
@@ -33,10 +44,9 @@ namespace CodexTests.BasicTests
|
||||
.AsStorageNode()
|
||||
.AsValidator()));
|
||||
|
||||
var expectedHostBalance = (numberOfHosts * hostInitialBalance.TstWei).TstWei();
|
||||
foreach (var host in hosts)
|
||||
{
|
||||
AssertBalance(contracts, host, Is.EqualTo(expectedHostBalance));
|
||||
AssertBalance(contracts, host, Is.EqualTo(hostInitialBalance));
|
||||
|
||||
var availability = new StorageAvailability(
|
||||
totalSpace: 10.GB(),
|
||||
@@ -47,7 +57,7 @@ namespace CodexTests.BasicTests
|
||||
host.Marketplace.MakeStorageAvailable(availability);
|
||||
}
|
||||
|
||||
var testFile = GenerateTestFile(fileSize);
|
||||
var testFile = CreateFile(fileSize);
|
||||
|
||||
var client = StartCodex(s => s
|
||||
.WithName("Client")
|
||||
@@ -56,20 +66,32 @@ namespace CodexTests.BasicTests
|
||||
|
||||
AssertBalance(contracts, client, Is.EqualTo(clientInitialBalance));
|
||||
|
||||
var contentId = client.UploadFile(testFile);
|
||||
var uploadCid = client.UploadFile(testFile);
|
||||
|
||||
var purchase = new StoragePurchaseRequest(contentId)
|
||||
var purchase = new StoragePurchaseRequest(uploadCid)
|
||||
{
|
||||
PricePerSlotPerSecond = 2.TstWei(),
|
||||
RequiredCollateral = 10.TstWei(),
|
||||
MinRequiredNumberOfNodes = 5,
|
||||
NodeFailureTolerance = 2,
|
||||
ProofProbability = 5,
|
||||
Duration = TimeSpan.FromMinutes(5),
|
||||
Expiry = TimeSpan.FromMinutes(4)
|
||||
Duration = TimeSpan.FromMinutes(6),
|
||||
Expiry = TimeSpan.FromMinutes(5)
|
||||
};
|
||||
|
||||
var purchaseContract = client.Marketplace.RequestStorage(purchase);
|
||||
|
||||
var contractCid = purchaseContract.ContentId;
|
||||
Assert.That(uploadCid.Id, Is.Not.EqualTo(contractCid.Id));
|
||||
|
||||
// Download both from client.
|
||||
testFile.AssertIsEqual(client.DownloadContent(uploadCid));
|
||||
testFile.AssertIsEqual(client.DownloadContent(contractCid));
|
||||
|
||||
// Download both from another node.
|
||||
var downloader = StartCodex(s => s.WithName("Downloader"));
|
||||
testFile.AssertIsEqual(downloader.DownloadContent(uploadCid));
|
||||
testFile.AssertIsEqual(downloader.DownloadContent(contractCid));
|
||||
|
||||
WaitForAllSlotFilledEvents(contracts, purchase, geth);
|
||||
|
||||
@@ -86,12 +108,13 @@ namespace CodexTests.BasicTests
|
||||
}
|
||||
|
||||
[Test]
|
||||
[Ignore("Integrated into MarketplaceExample to speed up testing.")]
|
||||
public void CanDownloadContentFromContractCid()
|
||||
{
|
||||
var fileSize = 10.MB();
|
||||
var geth = Ci.StartGethNode(s => s.IsMiner().WithName("disttest-geth"));
|
||||
var contracts = Ci.StartCodexContracts(geth);
|
||||
var testFile = GenerateTestFile(fileSize);
|
||||
var testFile = CreateFile(fileSize);
|
||||
|
||||
var client = StartCodex(s => s
|
||||
.WithName("Client")
|
||||
@@ -125,12 +148,24 @@ namespace CodexTests.BasicTests
|
||||
testFile.AssertIsEqual(downloader.DownloadContent(contractCid));
|
||||
}
|
||||
|
||||
private TrackedFile CreateFile(ByteSize fileSize)
|
||||
{
|
||||
var segmentSize = new ByteSize(fileSize.SizeInBytes / 4);
|
||||
|
||||
return GenerateTestFile(o => o
|
||||
.Random(segmentSize)
|
||||
.ByteRepeat(new byte[] { 0xaa }, segmentSize)
|
||||
.Random(segmentSize)
|
||||
.ByteRepeat(new byte[] { 0xee }, segmentSize)
|
||||
);
|
||||
}
|
||||
|
||||
private void WaitForAllSlotFilledEvents(ICodexContracts contracts, StoragePurchaseRequest purchase, IGethNode geth)
|
||||
{
|
||||
Time.Retry(() =>
|
||||
{
|
||||
var blockRange = geth.ConvertTimeRangeToBlockRange(GetTestRunTimeRange());
|
||||
var slotFilledEvents = contracts.GetSlotFilledEvents(blockRange);
|
||||
var events = contracts.GetEvents(GetTestRunTimeRange());
|
||||
var slotFilledEvents = events.GetSlotFilledEvents();
|
||||
|
||||
var msg = $"SlotFilledEvents: {slotFilledEvents.Length} - NumSlots: {purchase.MinRequiredNumberOfNodes}";
|
||||
Debug(msg);
|
||||
@@ -147,7 +182,8 @@ namespace CodexTests.BasicTests
|
||||
|
||||
private Request GetOnChainStorageRequest(ICodexContracts contracts, IGethNode geth)
|
||||
{
|
||||
var requests = contracts.GetStorageRequests(geth.ConvertTimeRangeToBlockRange(GetTestRunTimeRange()));
|
||||
var events = contracts.GetEvents(GetTestRunTimeRange());
|
||||
var requests = events.GetStorageRequests();
|
||||
Assert.That(requests.Length, Is.EqualTo(1));
|
||||
return requests.Single();
|
||||
}
|
||||
|
||||
@@ -0,0 +1,88 @@
|
||||
using CodexPlugin;
|
||||
using NUnit.Framework;
|
||||
using System;
|
||||
using System.Collections.Generic;
|
||||
using System.Linq;
|
||||
using System.Text;
|
||||
using System.Threading.Tasks;
|
||||
using Utils;
|
||||
|
||||
namespace CodexTests.BasicTests
|
||||
{
|
||||
[TestFixture]
|
||||
public class PyramidTests : CodexDistTest
|
||||
{
|
||||
[Test]
|
||||
[CreateTranscript(nameof(PyramidTest))]
|
||||
public void PyramidTest()
|
||||
{
|
||||
var size = 5.MB();
|
||||
var numberOfLayers = 4;
|
||||
|
||||
var bottomLayer = StartLayers(numberOfLayers);
|
||||
|
||||
var cids = UploadFiles(bottomLayer, size);
|
||||
|
||||
DownloadAllFilesFromEachNodeInLayer(bottomLayer, cids);
|
||||
}
|
||||
|
||||
private List<ICodexNode> StartLayers(int numberOfLayers)
|
||||
{
|
||||
var layer = new List<ICodexNode>();
|
||||
layer.Add(StartCodex(s => s.WithName("Top")));
|
||||
|
||||
for (var i = 0; i < numberOfLayers; i++)
|
||||
{
|
||||
var newLayer = new List<ICodexNode>();
|
||||
foreach (var node in layer)
|
||||
{
|
||||
newLayer.AddRange(StartCodex(2, s => s.WithBootstrapNode(node).WithName("Layer[" + i + "]")));
|
||||
}
|
||||
|
||||
layer.Clear();
|
||||
layer.AddRange(newLayer);
|
||||
}
|
||||
|
||||
return layer;
|
||||
}
|
||||
|
||||
private ContentId[] UploadFiles(List<ICodexNode> layer, ByteSize size)
|
||||
{
|
||||
var uploadTasks = new List<Task<ContentId>>();
|
||||
foreach (var node in layer)
|
||||
{
|
||||
uploadTasks.Add(Task.Run<ContentId>(() =>
|
||||
{
|
||||
var file = GenerateTestFile(size);
|
||||
return node.UploadFile(file);
|
||||
}));
|
||||
}
|
||||
|
||||
var cids = uploadTasks.Select(t =>
|
||||
{
|
||||
t.Wait();
|
||||
return t.Result;
|
||||
}).ToArray();
|
||||
|
||||
return cids;
|
||||
}
|
||||
|
||||
private void DownloadAllFilesFromEachNodeInLayer(List<ICodexNode> layer, ContentId[] cids)
|
||||
{
|
||||
var downloadTasks = new List<Task>();
|
||||
foreach (var node in layer)
|
||||
{
|
||||
downloadTasks.Add(Task.Run(() =>
|
||||
{
|
||||
var dlCids = RandomUtils.Shuffled(cids);
|
||||
foreach (var cid in dlCids)
|
||||
{
|
||||
node.DownloadContent(cid);
|
||||
}
|
||||
}));
|
||||
}
|
||||
|
||||
Task.WaitAll(downloadTasks.ToArray());
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -22,6 +22,24 @@ namespace CodexTests.BasicTests
|
||||
testFile.AssertIsEqual(downloadedFile);
|
||||
}
|
||||
|
||||
[Test]
|
||||
[CreateTranscript(nameof(SwarmTest))]
|
||||
public void SwarmTest()
|
||||
{
|
||||
var uploader = StartCodex(s => s.WithName("uploader"));
|
||||
var downloaders = StartCodex(5, s => s.WithName("downloader"));
|
||||
|
||||
var file = GenerateTestFile(100.MB());
|
||||
var cid = uploader.UploadFile(file);
|
||||
|
||||
var result = Parallel.ForEach(downloaders, d =>
|
||||
{
|
||||
d.DownloadContent(cid);
|
||||
});
|
||||
|
||||
Assert.That(result.IsCompleted);
|
||||
}
|
||||
|
||||
[Test]
|
||||
public void DownloadingUnknownCidDoesNotCauseCrash()
|
||||
{
|
||||
|
||||
@@ -1,19 +1,25 @@
|
||||
using CodexContractsPlugin;
|
||||
using CodexNetDeployer;
|
||||
using CodexPlugin;
|
||||
using CodexPlugin.OverwatchSupport;
|
||||
using CodexTests.Helpers;
|
||||
using Core;
|
||||
using DistTestCore;
|
||||
using DistTestCore.Helpers;
|
||||
using DistTestCore.Logs;
|
||||
using Logging;
|
||||
using MetricsPlugin;
|
||||
using Newtonsoft.Json;
|
||||
using NUnit.Framework;
|
||||
using NUnit.Framework.Constraints;
|
||||
using OverwatchTranscript;
|
||||
|
||||
namespace CodexTests
|
||||
{
|
||||
public class CodexDistTest : DistTest
|
||||
{
|
||||
private static readonly Dictionary<TestLifecycle, CodexTranscriptWriter> writers = new Dictionary<TestLifecycle, CodexTranscriptWriter>();
|
||||
|
||||
public CodexDistTest()
|
||||
{
|
||||
ProjectPlugin.Load<CodexPlugin.CodexPlugin>();
|
||||
@@ -29,6 +35,18 @@ namespace CodexTests
|
||||
localBuilder.Build();
|
||||
}
|
||||
|
||||
protected override void LifecycleStart(TestLifecycle lifecycle)
|
||||
{
|
||||
base.LifecycleStart(lifecycle);
|
||||
SetupTranscript(lifecycle);
|
||||
}
|
||||
|
||||
protected override void LifecycleStop(TestLifecycle lifecycle, DistTestResult result)
|
||||
{
|
||||
base.LifecycleStop(lifecycle, result);
|
||||
TeardownTranscript(lifecycle, result);
|
||||
}
|
||||
|
||||
public ICodexNode StartCodex()
|
||||
{
|
||||
return StartCodex(s => { });
|
||||
@@ -108,5 +126,82 @@ namespace CodexTests
|
||||
protected virtual void OnCodexSetup(ICodexSetup setup)
|
||||
{
|
||||
}
|
||||
|
||||
private CreateTranscriptAttribute? GetTranscriptAttributeOfCurrentTest()
|
||||
{
|
||||
var attrs = GetCurrentTestMethodAttribute<CreateTranscriptAttribute>();
|
||||
if (attrs.Any()) return attrs.Single();
|
||||
return null;
|
||||
}
|
||||
|
||||
private void SetupTranscript(TestLifecycle lifecycle)
|
||||
{
|
||||
var attr = GetTranscriptAttributeOfCurrentTest();
|
||||
if (attr == null) return;
|
||||
|
||||
var config = new CodexTranscriptWriterConfig(
|
||||
attr.IncludeBlockReceivedEvents
|
||||
);
|
||||
|
||||
var log = new LogPrefixer(lifecycle.Log, "(Transcript) ");
|
||||
var writer = new CodexTranscriptWriter(log, config, Transcript.NewWriter(log));
|
||||
Ci.SetCodexHooksProvider(writer);
|
||||
writers.Add(lifecycle, writer);
|
||||
}
|
||||
|
||||
private void TeardownTranscript(TestLifecycle lifecycle, DistTestResult result)
|
||||
{
|
||||
var attr = GetTranscriptAttributeOfCurrentTest();
|
||||
if (attr == null) return;
|
||||
|
||||
var outputFilepath = GetOutputFullPath(lifecycle, attr);
|
||||
|
||||
var writer = writers[lifecycle];
|
||||
writers.Remove(lifecycle);
|
||||
|
||||
writer.AddResult(result.Success, result.Result);
|
||||
|
||||
try
|
||||
{
|
||||
Stopwatch.Measure(lifecycle.Log, "Transcript.ProcessLogs", () =>
|
||||
{
|
||||
writer.ProcessLogs(lifecycle.DownloadAllLogs());
|
||||
});
|
||||
|
||||
Stopwatch.Measure(lifecycle.Log, $"Transcript.Finalize: {outputFilepath}", () =>
|
||||
{
|
||||
writer.IncludeFile(lifecycle.Log.LogFile.FullFilename);
|
||||
writer.Finalize(outputFilepath);
|
||||
});
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
lifecycle.Log.Error("Failure during transcript teardown: " + ex);
|
||||
}
|
||||
}
|
||||
|
||||
private string GetOutputFullPath(TestLifecycle lifecycle, CreateTranscriptAttribute attr)
|
||||
{
|
||||
var outputPath = Path.GetDirectoryName(lifecycle.Log.LogFile.FullFilename);
|
||||
if (outputPath == null) throw new Exception("Logfile path is null");
|
||||
var filename = Path.GetFileNameWithoutExtension(lifecycle.Log.LogFile.FullFilename);
|
||||
if (string.IsNullOrEmpty(filename)) throw new Exception("Logfile name is null or empty");
|
||||
var outputFile = Path.Combine(outputPath, filename + "_" + attr.OutputFilename);
|
||||
if (!outputFile.EndsWith(".owts")) outputFile += ".owts";
|
||||
return outputFile;
|
||||
}
|
||||
}
|
||||
|
||||
[AttributeUsage(AttributeTargets.Method, AllowMultiple = false)]
|
||||
public class CreateTranscriptAttribute : PropertyAttribute
|
||||
{
|
||||
public CreateTranscriptAttribute(string outputFilename, bool includeBlockReceivedEvents = true)
|
||||
{
|
||||
OutputFilename = outputFilename;
|
||||
IncludeBlockReceivedEvents = includeBlockReceivedEvents;
|
||||
}
|
||||
|
||||
public string OutputFilename { get; }
|
||||
public bool IncludeBlockReceivedEvents { get; }
|
||||
}
|
||||
}
|
||||
|
||||
@@ -13,11 +13,13 @@
|
||||
</ItemGroup>
|
||||
|
||||
<ItemGroup>
|
||||
<ProjectReference Include="..\..\Framework\DiscordRewards\DiscordRewards.csproj" />
|
||||
<ProjectReference Include="..\..\ProjectPlugins\CodexContractsPlugin\CodexContractsPlugin.csproj" />
|
||||
<ProjectReference Include="..\..\ProjectPlugins\CodexDiscordBotPlugin\CodexDiscordBotPlugin.csproj" />
|
||||
<ProjectReference Include="..\..\ProjectPlugins\CodexPlugin\CodexPlugin.csproj" />
|
||||
<ProjectReference Include="..\..\ProjectPlugins\GethPlugin\GethPlugin.csproj" />
|
||||
<ProjectReference Include="..\..\ProjectPlugins\MetricsPlugin\MetricsPlugin.csproj" />
|
||||
<ProjectReference Include="..\..\Tools\TestNetRewarder\TestNetRewarder.csproj" />
|
||||
<ProjectReference Include="..\DistTestCore\DistTestCore.csproj" />
|
||||
</ItemGroup>
|
||||
|
||||
|
||||
@@ -49,6 +49,8 @@ namespace CodexTests.PeerDiscoveryTests
|
||||
|
||||
private void AssertAllNodesConnected(IEnumerable<ICodexNode> nodes)
|
||||
{
|
||||
nodes = nodes.Concat(new[] { BootstrapNode }).ToArray()!;
|
||||
|
||||
CreatePeerConnectionTestHelpers().AssertFullyConnected(nodes);
|
||||
CheckRoutingTable(nodes);
|
||||
}
|
||||
|
||||
@@ -1,7 +1,13 @@
|
||||
using CodexContractsPlugin;
|
||||
using CodexDiscordBotPlugin;
|
||||
using CodexPlugin;
|
||||
using Core;
|
||||
using DiscordRewards;
|
||||
using DistTestCore;
|
||||
using GethPlugin;
|
||||
using KubernetesWorkflow.Types;
|
||||
using Logging;
|
||||
using Newtonsoft.Json;
|
||||
using NUnit.Framework;
|
||||
using Utils;
|
||||
|
||||
@@ -10,21 +16,134 @@ namespace CodexTests.UtilityTests
|
||||
[TestFixture]
|
||||
public class DiscordBotTests : AutoBootstrapDistTest
|
||||
{
|
||||
private readonly RewardRepo repo = new RewardRepo();
|
||||
private readonly TestToken hostInitialBalance = 3000000.TstWei();
|
||||
private readonly TestToken clientInitialBalance = 1000000000.TstWei();
|
||||
private readonly EthAccount clientAccount = EthAccount.GenerateNew();
|
||||
private readonly List<EthAccount> hostAccounts = new List<EthAccount>();
|
||||
private readonly List<ulong> rewardsSeen = new List<ulong>();
|
||||
private readonly TimeSpan rewarderInterval = TimeSpan.FromMinutes(1);
|
||||
private readonly List<string> receivedEvents = new List<string>();
|
||||
|
||||
[Test]
|
||||
[Ignore("Used for debugging bots")]
|
||||
[DontDownloadLogs]
|
||||
[Ignore("Used to debug testnet bots.")]
|
||||
public void BotRewardTest()
|
||||
{
|
||||
var myAccount = EthAccount.GenerateNew();
|
||||
|
||||
var sellerInitialBalance = 234.TstWei();
|
||||
var buyerInitialBalance = 100000.TstWei();
|
||||
var fileSize = 11.MB();
|
||||
|
||||
var geth = Ci.StartGethNode(s => s.IsMiner().WithName("disttest-geth"));
|
||||
var contracts = Ci.StartCodexContracts(geth);
|
||||
var gethInfo = CreateGethInfo(geth, contracts);
|
||||
|
||||
// start bot and rewarder
|
||||
var gethInfo = new DiscordBotGethInfo(
|
||||
var botContainer = StartDiscordBot(gethInfo);
|
||||
var rewarderContainer = StartRewarderBot(gethInfo, botContainer);
|
||||
|
||||
StartHosts(geth, contracts);
|
||||
var client = StartClient(geth, contracts);
|
||||
|
||||
var apiCalls = new RewardApiCalls(GetTestLog(), Ci, botContainer);
|
||||
apiCalls.Start(OnCommand);
|
||||
|
||||
var purchaseContract = ClientPurchasesStorage(client);
|
||||
purchaseContract.WaitForStorageContractStarted();
|
||||
purchaseContract.WaitForStorageContractFinished();
|
||||
Thread.Sleep(rewarderInterval * 3);
|
||||
|
||||
apiCalls.Stop();
|
||||
|
||||
AssertEventOccurance("Created as New.", 1);
|
||||
AssertEventOccurance("SlotFilled", Convert.ToInt32(GetNumberOfRequiredHosts()));
|
||||
AssertEventOccurance("Transit: New -> Started", 1);
|
||||
AssertEventOccurance("Transit: Started -> Finished", 1);
|
||||
|
||||
foreach (var r in repo.Rewards)
|
||||
{
|
||||
var seen = rewardsSeen.Any(s => r.RoleId == s);
|
||||
|
||||
Log($"{Lookup(r.RoleId)} = {seen}");
|
||||
}
|
||||
|
||||
Assert.That(repo.Rewards.All(r => rewardsSeen.Contains(r.RoleId)));
|
||||
}
|
||||
|
||||
private string Lookup(ulong rewardId)
|
||||
{
|
||||
var reward = repo.Rewards.Single(r => r.RoleId == rewardId);
|
||||
return $"({rewardId})'{reward.Message}'";
|
||||
}
|
||||
|
||||
private void AssertEventOccurance(string msg, int expectedCount)
|
||||
{
|
||||
Assert.That(receivedEvents.Count(e => e.Contains(msg)), Is.EqualTo(expectedCount),
|
||||
$"Event '{msg}' did not occure correct number of times.");
|
||||
}
|
||||
|
||||
private void OnCommand(string timestamp, GiveRewardsCommand call)
|
||||
{
|
||||
Log($"<API call {timestamp}>");
|
||||
receivedEvents.AddRange(call.EventsOverview);
|
||||
foreach (var e in call.EventsOverview)
|
||||
{
|
||||
Log("\tEvent: " + e);
|
||||
}
|
||||
foreach (var r in call.Rewards)
|
||||
{
|
||||
var reward = repo.Rewards.Single(a => a.RoleId == r.RewardId);
|
||||
if (r.UserAddresses.Any()) rewardsSeen.Add(reward.RoleId);
|
||||
foreach (var address in r.UserAddresses)
|
||||
{
|
||||
var user = IdentifyAccount(address);
|
||||
Log("\tReward: " + user + ": " + reward.Message);
|
||||
}
|
||||
}
|
||||
Log($"</API call>");
|
||||
}
|
||||
|
||||
private IStoragePurchaseContract ClientPurchasesStorage(ICodexNode client)
|
||||
{
|
||||
var testFile = GenerateTestFile(GetMinFileSize());
|
||||
var contentId = client.UploadFile(testFile);
|
||||
var purchase = new StoragePurchaseRequest(contentId)
|
||||
{
|
||||
PricePerSlotPerSecond = 2.TstWei(),
|
||||
RequiredCollateral = 10.TstWei(),
|
||||
MinRequiredNumberOfNodes = GetNumberOfRequiredHosts(),
|
||||
NodeFailureTolerance = 2,
|
||||
ProofProbability = 5,
|
||||
Duration = GetMinRequiredRequestDuration(),
|
||||
Expiry = GetMinRequiredRequestDuration() - TimeSpan.FromMinutes(1)
|
||||
};
|
||||
|
||||
return client.Marketplace.RequestStorage(purchase);
|
||||
}
|
||||
|
||||
private ICodexNode StartClient(IGethNode geth, ICodexContracts contracts)
|
||||
{
|
||||
var node = StartCodex(s => s
|
||||
.WithName("Client")
|
||||
.EnableMarketplace(geth, contracts, m => m
|
||||
.WithAccount(clientAccount)
|
||||
.WithInitial(10.Eth(), clientInitialBalance)));
|
||||
|
||||
Log($"Client {node.EthAccount.EthAddress}");
|
||||
return node;
|
||||
}
|
||||
|
||||
private RunningPod StartRewarderBot(DiscordBotGethInfo gethInfo, RunningContainer botContainer)
|
||||
{
|
||||
return Ci.DeployRewarderBot(new RewarderBotStartupConfig(
|
||||
name: "rewarder-bot",
|
||||
discordBotHost: botContainer.GetInternalAddress(DiscordBotContainerRecipe.RewardsPort).Host,
|
||||
discordBotPort: botContainer.GetInternalAddress(DiscordBotContainerRecipe.RewardsPort).Port,
|
||||
intervalMinutes: Convert.ToInt32(Math.Round(rewarderInterval.TotalMinutes)),
|
||||
historyStartUtc: DateTime.UtcNow,
|
||||
gethInfo: gethInfo,
|
||||
dataPath: null
|
||||
));
|
||||
}
|
||||
|
||||
private DiscordBotGethInfo CreateGethInfo(IGethNode geth, ICodexContracts contracts)
|
||||
{
|
||||
return new DiscordBotGethInfo(
|
||||
host: geth.Container.GetInternalAddress(GethContainerRecipe.HttpPortTag).Host,
|
||||
port: geth.Container.GetInternalAddress(GethContainerRecipe.HttpPortTag).Port,
|
||||
privKey: geth.StartResult.Account.PrivateKey,
|
||||
@@ -32,8 +151,12 @@ namespace CodexTests.UtilityTests
|
||||
tokenAddress: contracts.Deployment.TokenAddress,
|
||||
abi: contracts.Deployment.Abi
|
||||
);
|
||||
}
|
||||
|
||||
private RunningContainer StartDiscordBot(DiscordBotGethInfo gethInfo)
|
||||
{
|
||||
var bot = Ci.DeployCodexDiscordBot(new DiscordBotStartupConfig(
|
||||
name: "bot",
|
||||
name: "discord-bot",
|
||||
token: "aaa",
|
||||
serverName: "ThatBen's server",
|
||||
adminRoleName: "bottest-admins",
|
||||
@@ -42,70 +165,190 @@ namespace CodexTests.UtilityTests
|
||||
kubeNamespace: "notneeded",
|
||||
gethInfo: gethInfo
|
||||
));
|
||||
var botContainer = bot.Containers.Single();
|
||||
Ci.DeployRewarderBot(new RewarderBotStartupConfig(
|
||||
//discordBotHost: "http://" + botContainer.GetAddress(GetTestLog(), DiscordBotContainerRecipe.RewardsPort).Host,
|
||||
//discordBotPort: botContainer.GetAddress(GetTestLog(), DiscordBotContainerRecipe.RewardsPort).Port,
|
||||
discordBotHost: botContainer.GetInternalAddress(DiscordBotContainerRecipe.RewardsPort).Host,
|
||||
discordBotPort: botContainer.GetInternalAddress(DiscordBotContainerRecipe.RewardsPort).Port,
|
||||
intervalMinutes: "1",
|
||||
historyStartUtc: GetTestRunTimeRange().From - TimeSpan.FromMinutes(3),
|
||||
gethInfo: gethInfo,
|
||||
dataPath: null
|
||||
));
|
||||
return bot.Containers.Single();
|
||||
}
|
||||
|
||||
var numberOfHosts = 3;
|
||||
private void StartHosts(IGethNode geth, ICodexContracts contracts)
|
||||
{
|
||||
var hosts = StartCodex(GetNumberOfLiveHosts(), s => s
|
||||
.WithName("Host")
|
||||
.WithLogLevel(CodexLogLevel.Trace, new CodexLogCustomTopics(CodexLogLevel.Error, CodexLogLevel.Error, CodexLogLevel.Warn)
|
||||
{
|
||||
ContractClock = CodexLogLevel.Trace,
|
||||
})
|
||||
.WithStorageQuota(Mult(GetMinFileSizePlus(50), GetNumberOfLiveHosts()))
|
||||
.EnableMarketplace(geth, contracts, m => m
|
||||
.WithInitial(10.Eth(), hostInitialBalance)
|
||||
.AsStorageNode()
|
||||
.AsValidator()));
|
||||
|
||||
for (var i = 0; i < numberOfHosts; i++)
|
||||
var availability = new StorageAvailability(
|
||||
totalSpace: Mult(GetMinFileSize(), GetNumberOfLiveHosts()),
|
||||
maxDuration: TimeSpan.FromMinutes(30),
|
||||
minPriceForTotalSpace: 1.TstWei(),
|
||||
maxCollateral: hostInitialBalance
|
||||
);
|
||||
|
||||
foreach (var host in hosts)
|
||||
{
|
||||
var seller = StartCodex(s => s
|
||||
.WithName("Seller")
|
||||
.WithLogLevel(CodexLogLevel.Trace, new CodexLogCustomTopics(CodexLogLevel.Error, CodexLogLevel.Error, CodexLogLevel.Warn)
|
||||
{
|
||||
ContractClock = CodexLogLevel.Trace,
|
||||
})
|
||||
.WithStorageQuota(11.GB())
|
||||
.EnableMarketplace(geth, contracts, m => m
|
||||
.WithAccount(myAccount)
|
||||
.WithInitial(10.Eth(), sellerInitialBalance)
|
||||
.AsStorageNode()
|
||||
.AsValidator()));
|
||||
hostAccounts.Add(host.EthAccount);
|
||||
host.Marketplace.MakeStorageAvailable(availability);
|
||||
}
|
||||
}
|
||||
|
||||
var availability = new StorageAvailability(
|
||||
totalSpace: 10.GB(),
|
||||
maxDuration: TimeSpan.FromMinutes(30),
|
||||
minPriceForTotalSpace: 1.TstWei(),
|
||||
maxCollateral: 20.TstWei()
|
||||
);
|
||||
seller.Marketplace.MakeStorageAvailable(availability);
|
||||
private int GetNumberOfLiveHosts()
|
||||
{
|
||||
return Convert.ToInt32(GetNumberOfRequiredHosts()) + 3;
|
||||
}
|
||||
|
||||
private ByteSize Mult(ByteSize size, int mult)
|
||||
{
|
||||
return new ByteSize(size.SizeInBytes * mult);
|
||||
}
|
||||
|
||||
private ByteSize GetMinFileSizePlus(int plusMb)
|
||||
{
|
||||
return new ByteSize(GetMinFileSize().SizeInBytes + plusMb.MB().SizeInBytes);
|
||||
}
|
||||
|
||||
private ByteSize GetMinFileSize()
|
||||
{
|
||||
ulong minSlotSize = 0;
|
||||
ulong minNumHosts = 0;
|
||||
foreach (var r in repo.Rewards)
|
||||
{
|
||||
var s = Convert.ToUInt64(r.CheckConfig.MinSlotSize.SizeInBytes);
|
||||
var h = r.CheckConfig.MinNumberOfHosts;
|
||||
if (s > minSlotSize) minSlotSize = s;
|
||||
if (h > minNumHosts) minNumHosts = h;
|
||||
}
|
||||
|
||||
var testFile = GenerateTestFile(fileSize);
|
||||
var minFileSize = ((minSlotSize + 1024) * minNumHosts);
|
||||
return new ByteSize(Convert.ToInt64(minFileSize));
|
||||
}
|
||||
|
||||
var buyer = StartCodex(s => s
|
||||
.WithName("Buyer")
|
||||
.EnableMarketplace(geth, contracts, m => m
|
||||
.WithAccount(myAccount)
|
||||
.WithInitial(10.Eth(), buyerInitialBalance)));
|
||||
private uint GetNumberOfRequiredHosts()
|
||||
{
|
||||
return Convert.ToUInt32(repo.Rewards.Max(r => r.CheckConfig.MinNumberOfHosts));
|
||||
}
|
||||
|
||||
var contentId = buyer.UploadFile(testFile);
|
||||
private TimeSpan GetMinRequiredRequestDuration()
|
||||
{
|
||||
return repo.Rewards.Max(r => r.CheckConfig.MinDuration) + TimeSpan.FromSeconds(10);
|
||||
}
|
||||
|
||||
var purchase = new StoragePurchaseRequest(contentId)
|
||||
private string IdentifyAccount(string address)
|
||||
{
|
||||
if (address == clientAccount.EthAddress.Address) return "Client";
|
||||
try
|
||||
{
|
||||
PricePerSlotPerSecond = 2.TstWei(),
|
||||
RequiredCollateral = 10.TstWei(),
|
||||
MinRequiredNumberOfNodes = 5,
|
||||
NodeFailureTolerance = 2,
|
||||
ProofProbability = 5,
|
||||
Duration = TimeSpan.FromMinutes(6),
|
||||
Expiry = TimeSpan.FromMinutes(5)
|
||||
};
|
||||
var index = hostAccounts.FindIndex(a => a.EthAddress.Address == address);
|
||||
return "Host" + index;
|
||||
}
|
||||
catch
|
||||
{
|
||||
return "UNKNOWN";
|
||||
}
|
||||
}
|
||||
|
||||
var purchaseContract = buyer.Marketplace.RequestStorage(purchase);
|
||||
public class RewardApiCalls
|
||||
{
|
||||
private readonly ContainerFileMonitor monitor;
|
||||
|
||||
purchaseContract.WaitForStorageContractStarted();
|
||||
public RewardApiCalls(ILog log, CoreInterface ci, RunningContainer botContainer)
|
||||
{
|
||||
monitor = new ContainerFileMonitor(log, ci, botContainer, "/app/datapath/logs/discordbot.log");
|
||||
}
|
||||
|
||||
purchaseContract.WaitForStorageContractFinished();
|
||||
public void Start(Action<string, GiveRewardsCommand> onCommand)
|
||||
{
|
||||
monitor.Start(line => ParseLine(line, onCommand));
|
||||
}
|
||||
|
||||
public void Stop()
|
||||
{
|
||||
monitor.Stop();
|
||||
}
|
||||
|
||||
private void ParseLine(string line, Action<string, GiveRewardsCommand> onCommand)
|
||||
{
|
||||
try
|
||||
{
|
||||
var timestamp = line.Substring(0, 30);
|
||||
var json = line.Substring(31);
|
||||
|
||||
var cmd = JsonConvert.DeserializeObject<GiveRewardsCommand>(json);
|
||||
if (cmd != null)
|
||||
{
|
||||
onCommand(timestamp, cmd);
|
||||
}
|
||||
}
|
||||
catch
|
||||
{
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
public class ContainerFileMonitor
|
||||
{
|
||||
private readonly ILog log;
|
||||
private readonly CoreInterface ci;
|
||||
private readonly RunningContainer botContainer;
|
||||
private readonly string filePath;
|
||||
private readonly CancellationTokenSource cts = new CancellationTokenSource();
|
||||
private readonly List<string> seenLines = new List<string>();
|
||||
private Task worker = Task.CompletedTask;
|
||||
private Action<string> onNewLine = c => { };
|
||||
|
||||
public ContainerFileMonitor(ILog log, CoreInterface ci, RunningContainer botContainer, string filePath)
|
||||
{
|
||||
this.log = log;
|
||||
this.ci = ci;
|
||||
this.botContainer = botContainer;
|
||||
this.filePath = filePath;
|
||||
}
|
||||
|
||||
public void Start(Action<string> onNewLine)
|
||||
{
|
||||
this.onNewLine = onNewLine;
|
||||
worker = Task.Run(Worker);
|
||||
}
|
||||
|
||||
public void Stop()
|
||||
{
|
||||
cts.Cancel();
|
||||
worker.Wait();
|
||||
}
|
||||
|
||||
// did any container crash? that's why it repeats?
|
||||
|
||||
|
||||
private void Worker()
|
||||
{
|
||||
while (!cts.IsCancellationRequested)
|
||||
{
|
||||
Update();
|
||||
}
|
||||
}
|
||||
|
||||
private void Update()
|
||||
{
|
||||
Thread.Sleep(TimeSpan.FromSeconds(10));
|
||||
if (cts.IsCancellationRequested) return;
|
||||
|
||||
var botLog = ci.ExecuteContainerCommand(botContainer, "cat", filePath);
|
||||
var lines = botLog.Split(Environment.NewLine, StringSplitOptions.RemoveEmptyEntries);
|
||||
foreach (var line in lines)
|
||||
{
|
||||
// log.Log("line: " + line);
|
||||
|
||||
if (!seenLines.Contains(line))
|
||||
{
|
||||
seenLines.Add(line);
|
||||
onNewLine(line);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Some files were not shown because too many files have changed in this diff Show More
Reference in New Issue
Block a user