Compare commits

..
1 Commits
Author SHA1 Message Date
benbierens bb2f594138 investigate wip 2024-12-09 08:53:56 +01:00
231 changed files with 2184 additions and 5070 deletions
+4 -5
View File
@@ -1,6 +1,5 @@
name: Docker - AutoClient
on:
push:
branches:
@@ -13,15 +12,15 @@ on:
- 'Framework/**'
- 'ProjectPlugins/**'
- .github/workflows/docker-autoclient.yml
- .github/workflows/docker-reusable.yml
workflow_dispatch:
jobs:
build-and-push:
name: Build and Push
uses: codex-storage/github-actions/.github/workflows/docker-reusable.yml@master
uses: ./.github/workflows/docker-reusable.yml
with:
docker_file: Tools/AutoClient/docker/Dockerfile
dockerhub_repo: codexstorage/codex-autoclient
tag_latest: ${{ github.ref_name == github.event.repository.default_branch || startsWith(github.ref, 'refs/tags/') }}
docker_repo: codexstorage/codex-autoclient
secrets: inherit
+4 -3
View File
@@ -13,15 +13,16 @@ on:
- 'Framework/**'
- 'ProjectPlugins/**'
- .github/workflows/docker-discordbot.yml
- .github/workflows/docker-reusable.yml
workflow_dispatch:
jobs:
build-and-push:
name: Build and Push
uses: codex-storage/github-actions/.github/workflows/docker-reusable.yml@master
uses: ./.github/workflows/docker-reusable.yml
with:
docker_file: Tools/BiblioTech/docker/Dockerfile
dockerhub_repo: codexstorage/codex-discordbot
tag_latest: ${{ github.ref_name == github.event.repository.default_branch || startsWith(github.ref, 'refs/tags/') }}
docker_repo: codexstorage/codex-discordbot
secrets: inherit
+5 -4
View File
@@ -11,16 +11,17 @@ on:
- 'Tools/KeyMaker/**'
- 'Framework/**'
- 'ProjectPlugins/**'
- .github/workflows/docker-keymaker.yml
- .github/workflows/docker-KeyMaker.yml
- .github/workflows/docker-reusable.yml
workflow_dispatch:
jobs:
build-and-push:
name: Build and Push
uses: codex-storage/github-actions/.github/workflows/docker-reusable.yml@master
uses: ./.github/workflows/docker-reusable.yml
with:
docker_file: Tools/KeyMaker/docker/Dockerfile
dockerhub_repo: codexstorage/codex-keymaker
tag_latest: ${{ github.ref_name == github.event.repository.default_branch || startsWith(github.ref, 'refs/tags/') }}
docker_repo: codexstorage/codex-keymaker
secrets: inherit
+4 -3
View File
@@ -12,15 +12,16 @@ on:
- 'Framework/**'
- 'ProjectPlugins/**'
- .github/workflows/docker-marketinsights.yml
- .github/workflows/docker-reusable.yml
workflow_dispatch:
jobs:
build-and-push:
name: Build and Push
uses: codex-storage/github-actions/.github/workflows/docker-reusable.yml@master
uses: ./.github/workflows/docker-reusable.yml
with:
docker_file: Tools/MarketInsights/Dockerfile
dockerhub_repo: codexstorage/codex-marketinsights
tag_latest: ${{ github.ref_name == github.event.repository.default_branch || startsWith(github.ref, 'refs/tags/') }}
docker_repo: codexstorage/codex-marketinsights
secrets: inherit
+178
View File
@@ -0,0 +1,178 @@
name: Reusable - Docker
on:
workflow_call:
inputs:
docker_file:
default: docker/Dockerfile
description: Dockerfile
required: false
type: string
docker_repo:
default: codexstorage/cs-codex-dist-tests
description: DockerHub repository
required: false
type: string
tag_latest:
default: true
description: Set latest tag for Docker images
required: false
type: boolean
tag_sha:
default: true
description: Set Git short commit as Docker tag
required: false
type: boolean
tag_suffix:
default: ''
description: Suffix for Docker images tag
required: false
type: string
env:
DOCKER_FILE: ${{ inputs.docker_file }}
DOCKER_REPO: ${{ inputs.docker_repo }}
TAG_LATEST: ${{ inputs.tag_latest }}
TAG_SHA: ${{ inputs.tag_sha }}
TAG_SUFFIX: ${{ inputs.tag_suffix }}
jobs:
# Build platform specific image
build:
strategy:
fail-fast: true
matrix:
target:
- os: linux
arch: amd64
- os: linux
arch: arm64
include:
- target:
os: linux
arch: amd64
builder: ubuntu-22.04
- target:
os: linux
arch: arm64
builder: buildjet-4vcpu-ubuntu-2204-arm
name: Build ${{ matrix.target.os }}/${{ matrix.target.arch }}
runs-on: ${{ matrix.builder }}
env:
PLATFORM: ${{ format('{0}/{1}', 'linux', matrix.target.arch) }}
steps:
- name: Checkout
uses: actions/checkout@v4
- name: Docker - Meta
id: meta
uses: docker/metadata-action@v5
with:
images: ${{ env.DOCKER_REPO }}
- name: Docker - Set up Buildx
uses: docker/setup-buildx-action@v3
- name: Docker - Login to Docker Hub
uses: docker/login-action@v3
with:
username: ${{ secrets.DOCKERHUB_USERNAME }}
password: ${{ secrets.DOCKERHUB_TOKEN }}
- name: Docker - Build and Push by digest
id: build
uses: docker/build-push-action@v5
with:
context: .
file: ${{ env.DOCKER_FILE }}
platforms: ${{ env.PLATFORM }}
push: true
labels: ${{ steps.meta.outputs.labels }}
outputs: type=image,name=${{ env.DOCKER_REPO }},push-by-digest=true,name-canonical=true,push=true
- name: Docker - Export digest
run: |
mkdir -p /tmp/digests
digest="${{ steps.build.outputs.digest }}"
touch "/tmp/digests/${digest#sha256:}"
- name: Docker - Upload digest
uses: actions/upload-artifact@v4
with:
name: digests-${{ matrix.target.arch }}
path: /tmp/digests
if-no-files-found: error
retention-days: 1
# Publish multi-platform image
publish:
name: Publish multi-platform image
runs-on: ubuntu-latest
needs: build
steps:
- name: Docker - Variables
run: |
# Adjust custom suffix when set and
if [[ -n "${{ env.TAG_SUFFIX }}" ]]; then
echo "TAG_SUFFIX=-${{ env.TAG_SUFFIX }}" >>$GITHUB_ENV
fi
# Disable SHA tags on tagged release
if [[ ${{ startsWith(github.ref, 'refs/tags/') }} == "true" ]]; then
echo "TAG_SHA=false" >>$GITHUB_ENV
fi
# Handle latest and latest-custom using raw
if [[ ${{ env.TAG_SHA }} == "false" ]]; then
echo "TAG_LATEST=false" >>$GITHUB_ENV
echo "TAG_RAW=true" >>$GITHUB_ENV
if [[ -z "${{ env.TAG_SUFFIX }}" ]]; then
echo "TAG_RAW_VALUE=latest" >>$GITHUB_ENV
else
echo "TAG_RAW_VALUE=latest-{{ env.TAG_SUFFIX }}" >>$GITHUB_ENV
fi
else
echo "TAG_RAW=false" >>$GITHUB_ENV
fi
- name: Docker - Download digests
uses: actions/download-artifact@v4
with:
pattern: digests-*
merge-multiple: true
path: /tmp/digests
- name: Docker - Set up Buildx
uses: docker/setup-buildx-action@v3
- name: Docker - Meta
id: meta
uses: docker/metadata-action@v5
with:
images: ${{ env.DOCKER_REPO }}
flavor: |
latest=${{ env.TAG_LATEST }}
suffix=${{ env.TAG_SUFFIX }},onlatest=true
tags: |
type=semver,pattern={{version}}
type=raw,enable=${{ env.TAG_RAW }},value=latest
type=sha,enable=${{ env.TAG_SHA }}
- name: Docker - Login to Docker Hub
uses: docker/login-action@v3
with:
username: ${{ secrets.DOCKERHUB_USERNAME }}
password: ${{ secrets.DOCKERHUB_TOKEN }}
- name: Docker - Create manifest list and push
working-directory: /tmp/digests
run: |
docker buildx imagetools create $(jq -cr '.tags | map("-t " + .) | join(" ")' <<< "$DOCKER_METADATA_OUTPUT_JSON") \
$(printf '${{ env.DOCKER_REPO }}@sha256:%s ' *)
- name: Docker - Inspect image
run: |
docker buildx imagetools inspect ${{ env.DOCKER_REPO }}:${{ steps.meta.outputs.version }}
+4 -5
View File
@@ -1,6 +1,5 @@
name: Docker - Rewarder Bot
on:
push:
branches:
@@ -13,15 +12,15 @@ on:
- 'Framework/**'
- 'ProjectPlugins/**'
- .github/workflows/docker-rewarder.yml
- .github/workflows/docker-reusable.yml
workflow_dispatch:
jobs:
build-and-push:
name: Build and Push
uses: codex-storage/github-actions/.github/workflows/docker-reusable.yml@master
uses: ./.github/workflows/docker-reusable.yml
with:
docker_file: Tools/TestNetRewarder/docker/Dockerfile
dockerhub_repo: codexstorage/codex-rewarderbot
tag_latest: ${{ github.ref_name == github.event.repository.default_branch || startsWith(github.ref, 'refs/tags/') }}
docker_repo: codexstorage/codex-rewarderbot
secrets: inherit
+2 -5
View File
@@ -11,15 +11,12 @@ on:
- docker/Dockerfile
- docker/docker-entrypoint.sh
- .github/workflows/docker-runner.yml
- .github/workflows/docker-reusable.yml
workflow_dispatch:
jobs:
build-and-push:
name: Build and Push
uses: codex-storage/github-actions/.github/workflows/docker-reusable.yml@master
with:
docker_file: docker/Dockerfile
dockerhub_repo: codexstorage/cs-codex-dist-tests
tag_latest: ${{ github.ref_name == github.event.repository.default_branch || startsWith(github.ref, 'refs/tags/') }}
uses: ./.github/workflows/docker-reusable.yml
secrets: inherit
-1
View File
@@ -109,7 +109,6 @@ jobs:
# Get logs
while [[ $(kubectl get pod ${pod} -n ${namespace} -o jsonpath='{.status.phase}') == "Running" ]]; do
echo "Show ${pod} logs ..."
echo "----"
kubectl logs $pod -n $namespace -f || true
sleep 1
done
-1
View File
@@ -3,4 +3,3 @@ obj
bin
.vscode
Tools/AutoClient/datapath
.editorconfig
@@ -1,16 +0,0 @@
<Project Sdk="Microsoft.NET.Sdk">
<PropertyGroup>
<TargetFramework>net8.0</TargetFramework>
<ImplicitUsings>enable</ImplicitUsings>
<Nullable>enable</Nullable>
</PropertyGroup>
<ItemGroup>
<PackageReference Include="Nethereum.Web3" Version="4.14.0" />
</ItemGroup>
<ItemGroup>
<ProjectReference Include="..\Logging\Logging.csproj" />
</ItemGroup>
</Project>
-1
View File
@@ -9,7 +9,6 @@
<ItemGroup>
<ProjectReference Include="..\FileUtils\FileUtils.csproj" />
<ProjectReference Include="..\KubernetesWorkflow\KubernetesWorkflow.csproj" />
<ProjectReference Include="..\WebUtils\WebUtils.csproj" />
</ItemGroup>
</Project>
-1
View File
@@ -1,6 +1,5 @@
using KubernetesWorkflow;
using KubernetesWorkflow.Types;
using Logging;
namespace Core
{
@@ -1,10 +1,11 @@
using System.Net.Http.Headers;
using System.Net.Http.Json;
using Logging;
using Logging;
using Newtonsoft.Json;
using Serialization = Newtonsoft.Json.Serialization;
using System.Net.Http.Headers;
using System.Net.Http.Json;
using Utils;
namespace WebUtils
namespace Core
{
public interface IEndpoint
{
@@ -118,7 +119,7 @@ namespace WebUtils
var errors = new List<string>();
var deserialized = JsonConvert.DeserializeObject<T>(json, new JsonSerializerSettings()
{
Error = delegate (object? sender, Newtonsoft.Json.Serialization.ErrorEventArgs args)
Error = delegate (object? sender, Serialization.ErrorEventArgs args)
{
if (args.CurrentObject == args.ErrorContext.OriginalObject)
{
+3 -4
View File
@@ -1,6 +1,5 @@
using KubernetesWorkflow;
using Logging;
using WebUtils;
namespace Core
{
@@ -9,16 +8,16 @@ namespace Core
private readonly IToolsFactory toolsFactory;
private readonly PluginManager manager = new PluginManager();
public EntryPoint(ILog log, Configuration configuration, string fileManagerRootFolder, IWebCallTimeSet webCallTimeSet, IK8sTimeSet k8STimeSet)
public EntryPoint(ILog log, Configuration configuration, string fileManagerRootFolder, ITimeSet timeSet)
{
toolsFactory = new ToolsFactory(log, configuration, fileManagerRootFolder, webCallTimeSet, k8STimeSet);
toolsFactory = new ToolsFactory(log, configuration, fileManagerRootFolder, timeSet);
Tools = toolsFactory.CreateTools();
manager.InstantiatePlugins(PluginFinder.GetPluginTypes(), toolsFactory);
}
public EntryPoint(ILog log, Configuration configuration, string fileManagerRootFolder)
: this(log, configuration, fileManagerRootFolder, new DefaultWebCallTimeSet(), new DefaultK8sTimeSet())
: this(log, configuration, fileManagerRootFolder, new DefaultTimeSet())
{
}
@@ -1,7 +1,7 @@
using Logging;
using Utils;
namespace WebUtils
namespace Core
{
public interface IHttp
{
@@ -16,16 +16,16 @@ namespace WebUtils
private static object lockLock = new object();
private static readonly Dictionary<string, object> httpLocks = new Dictionary<string, object>();
private readonly ILog log;
private readonly IWebCallTimeSet timeSet;
private readonly ITimeSet timeSet;
private readonly Action<HttpClient> onClientCreated;
private readonly string id;
internal Http(string id, ILog log, IWebCallTimeSet timeSet)
internal Http(string id, ILog log, ITimeSet timeSet)
: this(id, log, timeSet, DoNothing)
{
}
internal Http(string id, ILog log, IWebCallTimeSet timeSet, Action<HttpClient> onClientCreated)
internal Http(string id, ILog log, ITimeSet timeSet, Action<HttpClient> onClientCreated)
{
this.id = id;
this.log = log;
+16 -15
View File
@@ -1,14 +1,12 @@
using FileUtils;
using KubernetesWorkflow;
using Logging;
using WebUtils;
namespace Core
{
public interface IPluginTools : IWorkflowTool, ILogTool, IHttpFactory, IFileTool
public interface IPluginTools : IWorkflowTool, ILogTool, IHttpFactoryTool, IFileTool
{
IWebCallTimeSet WebCallTimeSet { get; }
IK8sTimeSet K8STimeSet { get; }
ITimeSet TimeSet { get; }
/// <summary>
/// Deletes kubernetes and tracked file resources.
@@ -27,6 +25,13 @@ namespace Core
ILog GetLog();
}
public interface IHttpFactoryTool
{
IHttp CreateHttp(string id, Action<HttpClient> onClientCreated);
IHttp CreateHttp(string id, Action<HttpClient> onClientCreated, ITimeSet timeSet);
IHttp CreateHttp(string id);
}
public interface IFileTool
{
IFileManager GetFileManager();
@@ -35,22 +40,18 @@ namespace Core
internal class PluginTools : IPluginTools
{
private readonly WorkflowCreator workflowCreator;
private readonly HttpFactory httpFactory;
private readonly IFileManager fileManager;
private readonly LogPrefixer log;
internal PluginTools(ILog log, WorkflowCreator workflowCreator, string fileManagerRootFolder, IWebCallTimeSet webCallTimeSet, IK8sTimeSet k8STimeSet)
internal PluginTools(ILog log, WorkflowCreator workflowCreator, string fileManagerRootFolder, ITimeSet timeSet)
{
this.log = new LogPrefixer(log);
this.workflowCreator = workflowCreator;
httpFactory = new HttpFactory(log, webCallTimeSet);
WebCallTimeSet = webCallTimeSet;
K8STimeSet = k8STimeSet;
TimeSet = timeSet;
fileManager = new FileManager(log, fileManagerRootFolder);
}
public IWebCallTimeSet WebCallTimeSet { get; }
public IK8sTimeSet K8STimeSet { get; }
public ITimeSet TimeSet { get; }
public void ApplyLogPrefix(string prefix)
{
@@ -59,17 +60,17 @@ namespace Core
public IHttp CreateHttp(string id, Action<HttpClient> onClientCreated)
{
return httpFactory.CreateHttp(id, onClientCreated);
return CreateHttp(id, onClientCreated, TimeSet);
}
public IHttp CreateHttp(string id, Action<HttpClient> onClientCreated, IWebCallTimeSet timeSet)
public IHttp CreateHttp(string id, Action<HttpClient> onClientCreated, ITimeSet ts)
{
return httpFactory.CreateHttp(id, onClientCreated, timeSet);
return new Http(id, log, ts, onClientCreated);
}
public IHttp CreateHttp(string id)
{
return httpFactory.CreateHttp(id);
return new Http(id, log, TimeSet);
}
public IStartupWorkflow CreateWorkflow(string? namespaceOverride = null)
@@ -1,6 +1,6 @@
namespace WebUtils
namespace Core
{
public interface IWebCallTimeSet
public interface ITimeSet
{
/// <summary>
/// Timeout for a single HTTP call.
@@ -17,9 +17,20 @@
/// After a failed HTTP call, wait this long before trying again.
/// </summary>
TimeSpan HttpCallRetryDelay();
/// <summary>
/// After a failed K8s operation, wait this long before trying again.
/// </summary>
TimeSpan K8sOperationRetryDelay();
/// <summary>
/// Maximum total time to attempt to perform a successful k8s operation.
/// If k8s operations fail during this timespan, retries will be made.
/// </summary>
TimeSpan K8sOperationTimeout();
}
public class DefaultWebCallTimeSet : IWebCallTimeSet
public class DefaultTimeSet : ITimeSet
{
public TimeSpan HttpCallTimeout()
{
@@ -35,9 +46,19 @@
{
return TimeSpan.FromSeconds(1);
}
public TimeSpan K8sOperationRetryDelay()
{
return TimeSpan.FromSeconds(10);
}
public TimeSpan K8sOperationTimeout()
{
return TimeSpan.FromMinutes(30);
}
}
public class LongWebCallTimeSet : IWebCallTimeSet
public class LongTimeSet : ITimeSet
{
public TimeSpan HttpCallTimeout()
{
@@ -53,5 +74,15 @@
{
return TimeSpan.FromSeconds(20);
}
public TimeSpan K8sOperationRetryDelay()
{
return TimeSpan.FromSeconds(30);
}
public TimeSpan K8sOperationTimeout()
{
return TimeSpan.FromHours(1);
}
}
}
+4 -7
View File
@@ -1,6 +1,5 @@
using KubernetesWorkflow;
using Logging;
using WebUtils;
namespace Core
{
@@ -14,21 +13,19 @@ namespace Core
private readonly ILog log;
private readonly WorkflowCreator workflowCreator;
private readonly string fileManagerRootFolder;
private readonly IWebCallTimeSet webCallTimeSet;
private readonly IK8sTimeSet k8STimeSet;
private readonly ITimeSet timeSet;
public ToolsFactory(ILog log, Configuration configuration, string fileManagerRootFolder, IWebCallTimeSet webCallTimeSet, IK8sTimeSet k8STimeSet)
public ToolsFactory(ILog log, Configuration configuration, string fileManagerRootFolder, ITimeSet timeSet)
{
this.log = log;
workflowCreator = new WorkflowCreator(log, configuration);
this.fileManagerRootFolder = fileManagerRootFolder;
this.webCallTimeSet = webCallTimeSet;
this.k8STimeSet = k8STimeSet;
this.timeSet = timeSet;
}
public PluginTools CreateTools()
{
return new PluginTools(log, workflowCreator, fileManagerRootFolder, webCallTimeSet, k8STimeSet);
return new PluginTools(log, workflowCreator, fileManagerRootFolder, timeSet);
}
}
}
-6
View File
@@ -14,12 +14,6 @@ namespace FileUtils
Label = label;
}
public static TrackedFile FromPath(ILog log, string filepath)
{
// todo: I don't wanne have to do this to call upload.
return new TrackedFile(log, filepath, string.Empty);
}
public string Filename { get; }
public string Label { get; }
+2 -3
View File
@@ -1,5 +1,4 @@
using BlockchainUtils;
using CodexContractsPlugin;
using CodexContractsPlugin;
using CodexContractsPlugin.Marketplace;
using GethPlugin;
using Logging;
@@ -20,7 +19,7 @@ namespace GethConnector
return null;
}
var gethNode = new CustomGethNode(log, new BlockCache(), GethInput.GethHost, GethInput.GethPort, GethInput.PrivateKey);
var gethNode = new CustomGethNode(log, GethInput.GethHost, GethInput.GethPort, GethInput.PrivateKey);
var config = GetCodexMarketplaceConfig(gethNode, GethInput.MarketplaceAddress);
+2 -6
View File
@@ -45,12 +45,8 @@
private static string? GetEnvVar(List<string> error, string name)
{
var result = Environment.GetEnvironmentVariable(name);
if (string.IsNullOrEmpty(result))
{
error.Add($"'{name}' is not set.");
return null;
}
return result.Trim();
if (string.IsNullOrEmpty(result)) error.Add($"'{name}' is not set.");
return result;
}
}
}
@@ -3,7 +3,7 @@ using Logging;
namespace KubernetesWorkflow
{
public class ContainerCrashWatcher
public class CrashWatcher
{
private readonly ILog log;
private readonly KubernetesClientConfiguration config;
@@ -15,7 +15,7 @@ namespace KubernetesWorkflow
private Task? worker;
private Exception? workerException;
public ContainerCrashWatcher(ILog log, KubernetesClientConfiguration config, string containerName, string podName, string recipeName, string k8sNamespace)
public CrashWatcher(ILog log, KubernetesClientConfiguration config, string containerName, string podName, string recipeName, string k8sNamespace)
{
this.log = log;
this.config = config;
@@ -45,7 +45,7 @@ namespace KubernetesWorkflow
if (workerException != null) throw new Exception("Exception occurred in CrashWatcher worker thread.", workerException);
}
public bool HasCrashed()
public bool HasContainerCrashed()
{
using var client = new Kubernetes(config);
var result = HasContainerBeenRestarted(client);
@@ -70,13 +70,13 @@ namespace KubernetesWorkflow
using var client = new Kubernetes(config);
while (!token.IsCancellationRequested)
{
token.WaitHandle.WaitOne(TimeSpan.FromSeconds(10));
if (HasContainerBeenRestarted(client))
{
DownloadCrashedContainerLogs(client);
return;
}
token.WaitHandle.WaitOne(TimeSpan.FromSeconds(10));
}
}
@@ -91,7 +91,7 @@ namespace KubernetesWorkflow
private void DownloadCrashedContainerLogs(Kubernetes client)
{
using var stream = client.ReadNamespacedPodLog(podName, k8sNamespace, recipeName, previous: true);
var handler = new WriteToFileLogHandler(log, "Crash detected for " + containerName, containerName);
var handler = new WriteToFileLogHandler(log, "Crash detected for " + containerName);
handler.Log(stream);
}
}
@@ -1,8 +1,10 @@
namespace Logging
using Logging;
namespace KubernetesWorkflow
{
public interface IDownloadedLog
{
string SourceName { get; }
string ContainerName { get; }
void IterateLines(Action<string> action);
void IterateLines(Action<string> action, params string[] thatContain);
@@ -12,27 +14,21 @@
void DeleteFile();
}
public class DownloadedLog : IDownloadedLog
internal class DownloadedLog : IDownloadedLog
{
private readonly LogFile logFile;
public DownloadedLog(string filepath, string sourceName)
internal DownloadedLog(WriteToFileLogHandler logHandler, string containerName)
{
logFile = new LogFile(filepath);
SourceName = sourceName;
logFile = logHandler.LogFile;
ContainerName = containerName;
}
public DownloadedLog(LogFile logFile, string sourceName)
{
this.logFile = logFile;
SourceName = sourceName;
}
public string SourceName { get; }
public string ContainerName { get; }
public void IterateLines(Action<string> action)
{
using var file = File.OpenRead(logFile.Filename);
using var file = File.OpenRead(logFile.FullFilename);
using var streamReader = new StreamReader(file);
var line = streamReader.ReadLine();
@@ -68,12 +64,12 @@
public string GetFilepath()
{
return logFile.Filename;
return logFile.FullFilename;
}
public void DeleteFile()
{
File.Delete(logFile.Filename);
File.Delete(logFile.FullFilename);
}
}
}
+4 -17
View File
@@ -122,16 +122,6 @@ namespace KubernetesWorkflow
var all = client.Run(c => c.ListNamespace().Items);
var namespaces = all.Select(n => n.Name()).Where(n => n.StartsWith(prefix));
if (wait)
{
// If we're going to wait, trigger the delete for all the namespaces immediately.
// Then wait for them to finish one by one.
foreach (var ns in namespaces)
{
DeleteNamespace(ns, false);
}
}
foreach (var ns in namespaces)
{
DeleteNamespace(ns, wait);
@@ -745,10 +735,7 @@ namespace KubernetesWorkflow
throw new Exception($"Expected to find 1 pod by podLabel '{deployment.PodLabel}'. Found: {pods.Length}. " +
$"Total number of pods: {allPods.Items.Count}. Their labels: {string.Join(Environment.NewLine, allLabels)}");
}
var pod = pods[0];
if (pod.Status == null) throw new Exception("Pod status unknown");
if (string.IsNullOrEmpty(pod.Status.PodIP)) throw new Exception("Pod IP unknown");
return pod;
return pods[0];
}
#endregion
@@ -919,7 +906,7 @@ namespace KubernetesWorkflow
var msg = $"Pod crash detected for deployment {deploymentName} (pod:{podName})";
log.Error(msg);
DownloadPodLog(container, new WriteToFileLogHandler(log, msg, deploymentName), tailLines: null, previous: true);
DownloadPodLog(container, new WriteToFileLogHandler(log, msg), tailLines: null, previous: true);
throw new Exception(msg);
}
@@ -959,13 +946,13 @@ namespace KubernetesWorkflow
#endregion
public ContainerCrashWatcher CreateCrashWatcher(RunningContainer container)
public CrashWatcher CreateCrashWatcher(RunningContainer container)
{
var containerName = container.Name;
var podName = GetPodName(container);
var recipeName = container.Recipe.Name;
return new ContainerCrashWatcher(log, cluster.GetK8sClientConfig(), containerName, podName, recipeName, K8sNamespace);
return new CrashWatcher(log, cluster.GetK8sClientConfig(), containerName, podName, recipeName, K8sNamespace);
}
private V1Pod[] FindPodsByLabel(string podLabel)
@@ -1,42 +0,0 @@
namespace Core
{
public interface IK8sTimeSet
{
/// <summary>
/// After a failed K8s operation, wait this long before trying again.
/// </summary>
TimeSpan K8sOperationRetryDelay();
/// <summary>
/// Maximum total time to attempt to perform a successful k8s operation.
/// If k8s operations fail during this timespan, retries will be made.
/// </summary>
TimeSpan K8sOperationTimeout();
}
public class DefaultK8sTimeSet : IK8sTimeSet
{
public TimeSpan K8sOperationRetryDelay()
{
return TimeSpan.FromSeconds(10);
}
public TimeSpan K8sOperationTimeout()
{
return TimeSpan.FromMinutes(30);
}
}
public class LongK8sTimeSet : IK8sTimeSet
{
public TimeSpan K8sOperationRetryDelay()
{
return TimeSpan.FromSeconds(30);
}
public TimeSpan K8sOperationTimeout()
{
return TimeSpan.FromHours(1);
}
}
}
+3 -3
View File
@@ -25,11 +25,11 @@ namespace KubernetesWorkflow
public class WriteToFileLogHandler : LogHandler, ILogHandler
{
public WriteToFileLogHandler(ILog sourceLog, string description, string addFileName)
public WriteToFileLogHandler(ILog sourceLog, string description)
{
LogFile = sourceLog.CreateSubfile(addFileName);
LogFile = sourceLog.CreateSubfile();
var msg = $"{description} -->> {LogFile.Filename}";
var msg = $"{description} -->> {LogFile.FullFilename}";
sourceLog.Log(msg);
LogFile.Write(msg);
@@ -13,7 +13,7 @@ namespace KubernetesWorkflow
FutureContainers Start(int numberOfContainers, ILocation location, ContainerRecipeFactory recipeFactory, StartupConfig startupConfig);
PodInfo GetPodInfo(RunningContainer container);
PodInfo GetPodInfo(RunningPod pod);
ContainerCrashWatcher CreateCrashWatcher(RunningContainer container);
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);
@@ -61,8 +61,7 @@ namespace KubernetesWorkflow
var startResult = controller.BringOnline(recipes, location);
var containers = CreateContainers(startResult, recipes, startupConfig);
var info = GetPodInfo(startResult.Deployment);
var rc = new RunningPod(Guid.NewGuid().ToString(), info, startupConfig, startResult, containers);
var rc = new RunningPod(Guid.NewGuid().ToString(), startupConfig, startResult, containers);
cluster.Configuration.Hooks.OnContainersStarted(rc);
if (startResult.ExternalService != null)
@@ -84,14 +83,9 @@ namespace KubernetesWorkflow
});
}
public PodInfo GetPodInfo(RunningDeployment deployment)
{
return K8s(c => c.GetPodInfo(deployment));
}
public PodInfo GetPodInfo(RunningContainer container)
{
return GetPodInfo(container.RunningPod.StartResult.Deployment);
return K8s(c => c.GetPodInfo(container.RunningPod.StartResult.Deployment));
}
public PodInfo GetPodInfo(RunningPod pod)
@@ -99,7 +93,7 @@ namespace KubernetesWorkflow
return K8s(c => c.GetPodInfo(pod.StartResult.Deployment));
}
public ContainerCrashWatcher CreateCrashWatcher(RunningContainer container)
public CrashWatcher CreateCrashWatcher(RunningContainer container)
{
return K8s(c => c.CreateCrashWatcher(container));
}
@@ -133,14 +127,14 @@ namespace KubernetesWorkflow
{
var msg = $"Downloading container log for '{container.Name}'";
log.Log(msg);
var logHandler = new WriteToFileLogHandler(log, msg, container.Name);
var logHandler = new WriteToFileLogHandler(log, msg);
K8s(controller =>
{
controller.DownloadPodLog(container, logHandler, tailLines, previous);
});
return new DownloadedLog(logHandler.LogFile, container.Name);
return new DownloadedLog(logHandler, container.Name);
}
public string ExecuteCommand(RunningContainer container, string command, params string[] args)
@@ -215,7 +209,6 @@ namespace KubernetesWorkflow
var port = startResult.GetExternalServicePorts(recipe, tag);
return new Address(
logName: $"{recipe.Name}:{tag}",
startResult.Cluster.HostAddress,
port.Number);
}
@@ -227,7 +220,6 @@ namespace KubernetesWorkflow
var port = startResult.GetInternalServicePorts(recipe, tag);
return new Address(
logName: $"{serviceName}:{tag}",
$"http://{serviceName}.{namespaceName}.svc.cluster.local",
port.Number);
}
@@ -4,10 +4,9 @@ namespace KubernetesWorkflow.Types
{
public class RunningPod
{
public RunningPod(string id, PodInfo podInfo, StartupConfig startupConfig, StartResult startResult, RunningContainer[] containers)
public RunningPod(string id, StartupConfig startupConfig, StartResult startResult, RunningContainer[] containers)
{
Id = id;
PodInfo = podInfo;
StartupConfig = startupConfig;
StartResult = startResult;
Containers = containers;
@@ -16,7 +15,6 @@ namespace KubernetesWorkflow.Types
}
public string Id { get; }
public PodInfo PodInfo { get; }
public StartupConfig StartupConfig { get; }
public StartResult StartResult { get; }
public RunningContainer[] Containers { get; }
+5 -9
View File
@@ -8,7 +8,7 @@ namespace Logging
void Debug(string message = "", int skipFrames = 0);
void Error(string message);
void AddStringReplace(string from, string to);
LogFile CreateSubfile(string addName, string ext = "log");
LogFile CreateSubfile(string ext = "log");
}
public abstract class BaseLog : ILog
@@ -34,7 +34,7 @@ namespace Logging
{
get
{
if (logFile == null) logFile = new LogFile(GetFullName() + ".log");
if (logFile == null) logFile = new LogFile(GetFullName(), "log");
return logFile;
}
}
@@ -69,16 +69,12 @@ namespace Logging
public virtual void Delete()
{
File.Delete(LogFile.Filename);
File.Delete(LogFile.FullFilename);
}
public LogFile CreateSubfile(string addName, string ext = "log")
public LogFile CreateSubfile(string ext = "log")
{
addName = addName
.Replace("<", "")
.Replace(">", "");
return new LogFile($"{GetFullName()}_{GetSubfileNumber()}_{addName}.{ext}");
return new LogFile($"{GetFullName()}_{GetSubfileNumber()}", ext);
}
protected string ApplyReplacements(string str)
+17 -20
View File
@@ -1,19 +1,21 @@
using Utils;
namespace Logging
namespace Logging
{
public class LogFile
{
private readonly string extension;
private readonly object fileLock = new object();
private string filename;
public LogFile(string filename)
public LogFile(string filename, string extension)
{
Filename = filename;
this.filename = filename;
this.extension = extension;
FullFilename = filename + "." + extension;
EnsurePathExists(filename);
}
public string Filename { get; private set; }
public string FullFilename { get; private set; }
public void Write(string message)
{
@@ -26,7 +28,7 @@ namespace Logging
{
lock (fileLock)
{
File.AppendAllLines(Filename, new[] { message });
File.AppendAllLines(FullFilename, new[] { message });
}
}
catch (Exception ex)
@@ -35,24 +37,19 @@ namespace Logging
}
}
public void WriteRawMany(IEnumerable<string> lines)
public void ConcatToFilename(string toAdd)
{
try
{
lock (fileLock)
{
File.AppendAllLines(Filename, lines);
}
}
catch (Exception ex)
{
Console.WriteLine("Writing to log has failed: " + ex);
}
var oldFullName = FullFilename;
filename += toAdd;
FullFilename = filename + "." + extension;
File.Move(oldFullName, FullFilename);
}
private static string GetTimestamp()
{
return $"[{Time.FormatTimestamp(DateTime.UtcNow)}]";
return $"[{DateTime.UtcNow.ToString("o")}]";
}
private void EnsurePathExists(string filename)
+2 -2
View File
@@ -18,9 +18,9 @@
public string Prefix { get; set; } = string.Empty;
public LogFile CreateSubfile(string addName, string ext = "log")
public LogFile CreateSubfile(string ext = "log")
{
return backingLog.CreateSubfile(addName, ext);
return backingLog.CreateSubfile(ext);
}
public void Debug(string message = "", int skipFrames = 0)
+2 -2
View File
@@ -14,9 +14,9 @@
OnAll(l => l.AddStringReplace(from, to));
}
public LogFile CreateSubfile(string addName, string ext = "log")
public LogFile CreateSubfile(string ext = "log")
{
return targetLogs.First().CreateSubfile(addName, ext);
return targetLogs.First().CreateSubfile(ext);
}
public void Debug(string message = "", int skipFrames = 0)
@@ -1,4 +1,4 @@
namespace BlockchainUtils
namespace NethereumWorkflow.BlockUtils
{
public class BlockCache
{
@@ -1,4 +1,4 @@
namespace BlockchainUtils
namespace NethereumWorkflow.BlockUtils
{
public class BlockTimeEntry
{
@@ -1,6 +1,6 @@
using Logging;
namespace BlockchainUtils
namespace NethereumWorkflow.BlockUtils
{
public class BlockTimeFinder
{
@@ -1,11 +1,5 @@
namespace BlockchainUtils
namespace NethereumWorkflow.BlockUtils
{
public interface IWeb3Blocks
{
ulong GetCurrentBlockNumber();
DateTime? GetTimestampForBlock(ulong blockNumber);
}
public class BlockchainBounds
{
private readonly BlockCache cache;
@@ -1,7 +1,7 @@
using Nethereum.Hex.HexTypes;
using System.Numerics;
namespace BlockchainUtils
namespace NethereumWorkflow
{
public static class ConversionExtensions
{
@@ -1,25 +1,25 @@
using BlockchainUtils;
using Logging;
using Logging;
using Nethereum.ABI.FunctionEncoding.Attributes;
using Nethereum.Contracts;
using Nethereum.RPC.Eth.DTOs;
using Nethereum.Web3;
using NethereumWorkflow.BlockUtils;
using Utils;
namespace NethereumWorkflow
{
public class NethereumInteraction
{
private readonly BlockCache blockCache;
// BlockCache is a static instance: It stays alive for the duration of the application runtime.
private readonly static BlockCache blockCache = new BlockCache();
private readonly ILog log;
private readonly Web3 web3;
internal NethereumInteraction(ILog log, Web3 web3, BlockCache blockCache)
internal NethereumInteraction(ILog log, Web3 web3)
{
this.log = log;
this.web3 = web3;
this.blockCache = blockCache;
}
public string SendEth(string toAddress, decimal ethAmount)
@@ -50,21 +50,6 @@ namespace NethereumWorkflow
return Time.Wait(handler.QueryAsync<TResult>(contractAddress, function));
}
public TResult Call<TFunction, TResult>(string contractAddress, TFunction function, ulong blockNumber) where TFunction : FunctionMessage, new()
{
log.Debug(typeof(TFunction).ToString());
var handler = web3.Eth.GetContractQueryHandler<TFunction>();
return Time.Wait(handler.QueryAsync<TResult>(contractAddress, function, new BlockParameter(blockNumber)));
}
public void Call<TFunction>(string contractAddress, TFunction function, ulong blockNumber) where TFunction : FunctionMessage, new()
{
log.Debug(typeof(TFunction).ToString());
var handler = web3.Eth.GetContractQueryHandler<TFunction>();
var result = Time.Wait(handler.QueryRawAsync(contractAddress, function, new BlockParameter(blockNumber)));
var aaaa = 0;
}
public string SendTransaction<TFunction>(string contractAddress, TFunction function) where TFunction : FunctionMessage, new()
{
log.Debug();
@@ -1,5 +1,4 @@
using BlockchainUtils;
using Logging;
using Logging;
using Nethereum.Web3;
namespace NethereumWorkflow
@@ -7,15 +6,13 @@ namespace NethereumWorkflow
public class NethereumInteractionCreator
{
private readonly ILog log;
private readonly BlockCache blockCache;
private readonly string ip;
private readonly int port;
private readonly string privateKey;
public NethereumInteractionCreator(ILog log, BlockCache blockCache, string ip, int port, string privateKey)
public NethereumInteractionCreator(ILog log, string ip, int port, string privateKey)
{
this.log = log;
this.blockCache = blockCache;
this.ip = ip;
this.port = port;
this.privateKey = privateKey;
@@ -24,7 +21,7 @@ namespace NethereumWorkflow
public NethereumInteraction CreateWorkflow()
{
log.Debug("Starting interaction to " + ip + ":" + port);
return new NethereumInteraction(log, CreateWeb3(), blockCache);
return new NethereumInteraction(log, CreateWeb3());
}
private Web3 CreateWeb3()
@@ -12,7 +12,6 @@
</ItemGroup>
<ItemGroup>
<ProjectReference Include="..\BlockchainUtils\BlockchainUtils.csproj" />
<ProjectReference Include="..\Logging\Logging.csproj" />
<ProjectReference Include="..\Utils\Utils.csproj" />
</ItemGroup>
+7 -2
View File
@@ -1,11 +1,16 @@
using BlockchainUtils;
using Logging;
using Logging;
using Nethereum.RPC.Eth.DTOs;
using Nethereum.Web3;
using Utils;
namespace NethereumWorkflow
{
public interface IWeb3Blocks
{
ulong GetCurrentBlockNumber();
DateTime? GetTimestampForBlock(ulong blockNumber);
}
public class Web3Wrapper : IWeb3Blocks
{
private readonly Web3 web3;
+1 -19
View File
@@ -2,14 +2,12 @@
{
public class Address
{
public Address(string logName, string host, int port)
public Address(string host, int port)
{
LogName = logName;
Host = host;
Port = port;
}
public string LogName { get; }
public string Host { get; }
public int Port { get; }
@@ -22,21 +20,5 @@
{
return !string.IsNullOrEmpty(Host) && Port > 0;
}
public static Address Empty()
{
return new Address(string.Empty, string.Empty, 0);
}
}
public interface IHasMetricsScrapeTarget
{
Address GetMetricsScrapeTarget();
}
public interface IHasManyMetricScrapeTargets
{
Address[] GetMetricsScrapeTargets();
}
}
-64
View File
@@ -1,64 +0,0 @@
using System;
using System.Collections.Generic;
using System.Diagnostics.Contracts;
using System.Linq;
using System.Numerics;
using System.Text;
using System.Threading.Tasks;
namespace Utils
{
public static class Base58
{
private const string Digits = "123456789ABCDEFGHJKLMNPQRSTUVWXYZabcdefghijkmnopqrstuvwxyz";
public static string Encode(byte[] data)
{
// Decode byte[] to BigInteger
BigInteger intData = 0;
for (int i = 0; i < data.Length; i++)
{
intData = intData * 256 + data[i];
}
// Encode BigInteger to Base58 string
string result = "";
while (intData > 0)
{
int remainder = (int)(intData % 58);
intData /= 58;
result = Digits[remainder] + result;
}
// Append `1` for each leading 0 byte
for (int i = 0; i < data.Length && data[i] == 0; i++)
{
result = '1' + result;
}
return result;
}
public static byte[] Decode(string s)
{
BigInteger intData = 0;
for (int i = 0; i < s.Length; i++)
{
int digit = Digits.IndexOf(s[i]); //Slow
if (digit < 0)
throw new FormatException(string.Format("Invalid Base58 character `{0}` at position {1}", s[i], i));
intData = intData * 58 + digit;
}
// Encode BigInteger to byte[]
// Leading zero bytes get encoded as leading `1` characters
int leadingZeroCount = s.TakeWhile(c => c == '1').Count();
var leadingZeros = Enumerable.Repeat((byte)0, leadingZeroCount);
var bytesWithoutLeadingZeros =
intData.ToByteArray()
.Reverse()// to big endian
.SkipWhile(b => b == 0);//strip sign byte
var result = leadingZeros.Concat(bytesWithoutLeadingZeros).ToArray();
return result;
}
}
}
-9
View File
@@ -1,9 +0,0 @@
namespace Utils
{
//public interface ICrashWatcher
//{
// void Start();
// void Stop();
// bool HasCrashed();
//}
}
-20
View File
@@ -1,20 +0,0 @@
namespace Utils
{
[Serializable]
public class EthAccount
{
public EthAccount(EthAddress ethAddress, string privateKey)
{
EthAddress = ethAddress;
PrivateKey = privateKey;
}
public EthAddress EthAddress { get; }
public string PrivateKey { get; }
public override string ToString()
{
return EthAddress.ToString();
}
}
}
+1 -1
View File
@@ -98,7 +98,7 @@
private void Fail()
{
throw new TimeoutException($"Retry '{description}' timed out after {tryNumber} tries over {Time.FormatDuration(Duration())}: {GetFailureReport()}",
throw new TimeoutException($"Retry '{description}' timed out after {tryNumber} tries over {Time.FormatDuration(Duration())}: {GetFailureReport}",
new AggregateException(failures.Select(f => f.Exception)));
}
-5
View File
@@ -33,11 +33,6 @@
result += $"{d.Seconds} secs";
return result;
}
public static string FormatTimestamp(DateTime d)
{
return d.ToString("o");
}
public static TimeSpan ParseTimespan(string span)
{
-43
View File
@@ -1,43 +0,0 @@
using Logging;
namespace WebUtils
{
public interface IHttpFactory
{
IHttp CreateHttp(string id, Action<HttpClient> onClientCreated);
IHttp CreateHttp(string id, Action<HttpClient> onClientCreated, IWebCallTimeSet timeSet);
IHttp CreateHttp(string id);
}
public class HttpFactory : IHttpFactory
{
private readonly ILog log;
private readonly IWebCallTimeSet defaultTimeSet;
public HttpFactory(ILog log)
: this (log, new DefaultWebCallTimeSet())
{
}
public HttpFactory(ILog log, IWebCallTimeSet defaultTimeSet)
{
this.log = log;
this.defaultTimeSet = defaultTimeSet;
}
public IHttp CreateHttp(string id, Action<HttpClient> onClientCreated)
{
return CreateHttp(id, onClientCreated, defaultTimeSet);
}
public IHttp CreateHttp(string id, Action<HttpClient> onClientCreated, IWebCallTimeSet ts)
{
return new Http(id, log, ts, onClientCreated);
}
public IHttp CreateHttp(string id)
{
return new Http(id, log, defaultTimeSet);
}
}
}
-18
View File
@@ -1,18 +0,0 @@
<Project Sdk="Microsoft.NET.Sdk">
<PropertyGroup>
<TargetFramework>net8.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" />
<ProjectReference Include="..\Utils\Utils.csproj" />
</ItemGroup>
</Project>
@@ -1,35 +0,0 @@
<Project Sdk="Microsoft.NET.Sdk">
<PropertyGroup>
<TargetFramework>net8.0</TargetFramework>
<ImplicitUsings>enable</ImplicitUsings>
<Nullable>enable</Nullable>
</PropertyGroup>
<ItemGroup>
<None Remove="openapi.yaml" />
</ItemGroup>
<ItemGroup>
<OpenApiReference Include="openapi.yaml" CodeGenerator="NSwagCSharp" Namespace="CodexOpenApi" ClassName="CodexApiClient" />
</ItemGroup>
<ItemGroup>
<PackageReference Include="Microsoft.Extensions.ApiDescription.Client" Version="7.0.2">
<PrivateAssets>all</PrivateAssets>
<IncludeAssets>runtime; build; native; contentfiles; analyzers; buildtransitive</IncludeAssets>
</PackageReference>
<PackageReference Include="Newtonsoft.Json" Version="13.0.3" />
<PackageReference Include="NSwag.ApiDescription.Client" Version="13.18.2">
<PrivateAssets>all</PrivateAssets>
<IncludeAssets>runtime; build; native; contentfiles; analyzers; buildtransitive</IncludeAssets>
</PackageReference>
</ItemGroup>
<ItemGroup>
<ProjectReference Include="..\..\Framework\FileUtils\FileUtils.csproj" />
<ProjectReference Include="..\..\Framework\Logging\Logging.csproj" />
<ProjectReference Include="..\..\Framework\WebUtils\WebUtils.csproj" />
</ItemGroup>
</Project>
@@ -1,54 +0,0 @@
using Utils;
namespace CodexClient
{
public interface ICodexInstance
{
string Name { get; }
string ImageName { get; }
DateTime StartUtc { get; }
Address DiscoveryEndpoint { get; }
Address ApiEndpoint { get; }
Address ListenEndpoint { get; }
EthAccount? EthAccount { get; }
Address? MetricsEndpoint { get; }
}
public class CodexInstance : ICodexInstance
{
public CodexInstance(string name, string imageName, DateTime startUtc, Address discoveryEndpoint, Address apiEndpoint, Address listenEndpoint, EthAccount? ethAccount, Address? metricsEndpoint)
{
Name = name;
ImageName = imageName;
StartUtc = startUtc;
DiscoveryEndpoint = discoveryEndpoint;
ApiEndpoint = apiEndpoint;
ListenEndpoint = listenEndpoint;
EthAccount = ethAccount;
MetricsEndpoint = metricsEndpoint;
}
public string Name { get; }
public string ImageName { get; }
public DateTime StartUtc { get; }
public Address DiscoveryEndpoint { get; }
public Address ApiEndpoint { get; }
public Address ListenEndpoint { get; }
public EthAccount? EthAccount { get; }
public Address? MetricsEndpoint { get; }
public static ICodexInstance CreateFromApiEndpoint(string name, Address apiEndpoint, EthAccount? ethAccount = null)
{
return new CodexInstance(
name,
imageName: "-",
startUtc: DateTime.UtcNow,
discoveryEndpoint: Address.Empty(),
apiEndpoint: apiEndpoint,
listenEndpoint: Address.Empty(),
ethAccount: ethAccount,
metricsEndpoint: null
);
}
}
}
@@ -1,52 +0,0 @@
using CodexClient.Hooks;
using FileUtils;
using Logging;
using WebUtils;
namespace CodexClient
{
public class CodexNodeFactory
{
private readonly ILog log;
private readonly IFileManager fileManager;
private readonly CodexHooksFactory hooksFactory;
private readonly IHttpFactory httpFactory;
private readonly IProcessControlFactory processControlFactory;
public CodexNodeFactory(ILog log, IFileManager fileManager, CodexHooksFactory hooksFactory, IHttpFactory httpFactory, IProcessControlFactory processControlFactory)
{
this.log = log;
this.fileManager = fileManager;
this.hooksFactory = hooksFactory;
this.httpFactory = httpFactory;
this.processControlFactory = processControlFactory;
}
public CodexNodeFactory(ILog log, HttpFactory httpFactory, string dataDir)
: this(log, new FileManager(log, dataDir), new CodexHooksFactory(), httpFactory, new DoNothingProcessControlFactory())
{
}
public CodexNodeFactory(ILog log, string dataDir)
: this(log, new HttpFactory(log), dataDir)
{
}
public ICodexNode CreateCodexNode(ICodexInstance instance)
{
var processControl = processControlFactory.CreateProcessControl(instance);
var access = new CodexAccess(log, httpFactory, processControl, instance);
var hooks = hooksFactory.CreateHooks(access.GetName());
var marketplaceAccess = CreateMarketplaceAccess(instance, access, hooks);
var node = new CodexNode(log, access, fileManager, marketplaceAccess, hooks);
node.Initialize();
return node;
}
private IMarketplaceAccess CreateMarketplaceAccess(ICodexInstance instance, CodexAccess access, ICodexNodeHooks hooks)
{
if (instance.EthAccount == null) return new MarketplaceUnavailable();
return new MarketplaceAccess(log, access, hooks);
}
}
}
@@ -1,78 +0,0 @@
using Utils;
namespace CodexClient.Hooks
{
public interface ICodexNodeHooks
{
void OnNodeStarting(DateTime startUtc, string image, EthAccount? ethAccount);
void OnNodeStarted(ICodexNode node, 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);
}
public class MuxingCodexNodeHooks : ICodexNodeHooks
{
private readonly ICodexNodeHooks[] backingHooks;
public MuxingCodexNodeHooks(ICodexNodeHooks[] backingHooks)
{
this.backingHooks = backingHooks;
}
public void OnFileDownloaded(ByteSize size, ContentId cid)
{
foreach (var h in backingHooks) h.OnFileDownloaded(size, cid);
}
public void OnFileDownloading(ContentId cid)
{
foreach (var h in backingHooks) h.OnFileDownloading(cid);
}
public void OnFileUploaded(string uid, ByteSize size, ContentId cid)
{
foreach (var h in backingHooks) h.OnFileUploaded(uid, size, cid);
}
public void OnFileUploading(string uid, ByteSize size)
{
foreach (var h in backingHooks) h.OnFileUploading(uid, size);
}
public void OnNodeStarted(ICodexNode node, string peerId, string nodeId)
{
foreach (var h in backingHooks) h.OnNodeStarted(node, peerId, nodeId);
}
public void OnNodeStarting(DateTime startUtc, string image, EthAccount? ethAccount)
{
foreach (var h in backingHooks) h.OnNodeStarting(startUtc, image, ethAccount);
}
public void OnNodeStopping()
{
foreach (var h in backingHooks) h.OnNodeStopping();
}
public void OnStorageAvailabilityCreated(StorageAvailability response)
{
foreach (var h in backingHooks) h.OnStorageAvailabilityCreated(response);
}
public void OnStorageContractSubmitted(StoragePurchaseContract storagePurchaseContract)
{
foreach (var h in backingHooks) h.OnStorageContractSubmitted(storagePurchaseContract);
}
public void OnStorageContractUpdated(StoragePurchase purchaseStatus)
{
foreach (var h in backingHooks) h.OnStorageContractUpdated(purchaseStatus);
}
}
}
@@ -1,46 +0,0 @@
using Logging;
namespace CodexClient
{
public interface IProcessControlFactory
{
IProcessControl CreateProcessControl(ICodexInstance instance);
}
public interface IProcessControl
{
void Stop(bool waitTillStopped);
IDownloadedLog DownloadLog(LogFile file);
void DeleteDataDirFolder();
bool HasCrashed();
}
public class DoNothingProcessControlFactory : IProcessControlFactory
{
public IProcessControl CreateProcessControl(ICodexInstance instance)
{
return new DoNothingProcessControl();
}
}
public class DoNothingProcessControl : IProcessControl
{
public void DeleteDataDirFolder()
{
}
public IDownloadedLog DownloadLog(LogFile file)
{
throw new NotImplementedException("Not supported by DoNothingProcessControl");
}
public bool HasCrashed()
{
return false;
}
public void Stop(bool waitTillStopped)
{
}
}
}
@@ -14,8 +14,7 @@ namespace CodexContractsPlugin.ChainMonitor
RequestFailedEventDTO[] failed,
SlotFilledEventDTO[] slotFilled,
SlotFreedEventDTO[] slotFreed,
SlotReservationsFullEventDTO[] slotReservationsFull,
ProofSubmittedEventDTO[] proofSubmitted
SlotReservationsFullEventDTO[] slotReservationsFull
)
{
BlockInterval = blockInterval;
@@ -26,8 +25,8 @@ namespace CodexContractsPlugin.ChainMonitor
SlotFilled = slotFilled;
SlotFreed = slotFreed;
SlotReservationsFull = slotReservationsFull;
ProofSubmitted = proofSubmitted;
All = ConcatAll<IHasBlock>(requests, fulfilled, cancelled, failed, slotFilled, SlotFreed, SlotReservationsFull, ProofSubmitted);
All = ConcatAll<IHasBlock>(requests, fulfilled, cancelled, failed, slotFilled, SlotFreed, SlotReservationsFull);
}
public BlockInterval BlockInterval { get; }
@@ -38,7 +37,6 @@ namespace CodexContractsPlugin.ChainMonitor
public SlotFilledEventDTO[] SlotFilled { get; }
public SlotFreedEventDTO[] SlotFreed { get; }
public SlotReservationsFullEventDTO[] SlotReservationsFull { get; }
public ProofSubmittedEventDTO[] ProofSubmitted { get; }
public IHasBlock[] All { get; }
public static ChainEvents FromBlockInterval(ICodexContracts contracts, BlockInterval blockInterval)
@@ -61,8 +59,7 @@ namespace CodexContractsPlugin.ChainMonitor
events.GetRequestFailedEvents(),
events.GetSlotFilledEvents(),
events.GetSlotFreedEvents(),
events.GetSlotReservationsFullEvents(),
events.GetProofSubmittedEvents()
events.GetSlotReservationsFull()
);
}
@@ -1,7 +1,7 @@
using BlockchainUtils;
using CodexContractsPlugin.Marketplace;
using CodexContractsPlugin.Marketplace;
using GethPlugin;
using Logging;
using NethereumWorkflow.BlockUtils;
using System.Numerics;
using Utils;
@@ -17,7 +17,7 @@ namespace CodexContractsPlugin.ChainMonitor
void OnSlotFilled(RequestEvent requestEvent, EthAddress host, BigInteger slotIndex);
void OnSlotFreed(RequestEvent requestEvent, BigInteger slotIndex);
void OnSlotReservationsFull(RequestEvent requestEvent, BigInteger slotIndex);
void OnProofSubmitted(BlockTimeEntry block, string id);
void OnError(string msg);
}
@@ -39,21 +39,17 @@ namespace CodexContractsPlugin.ChainMonitor
private readonly ILog log;
private readonly ICodexContracts contracts;
private readonly IChainStateChangeHandler handler;
private readonly bool doProofPeriodMonitoring;
public ChainState(ILog log, ICodexContracts contracts, IChainStateChangeHandler changeHandler, DateTime startUtc, bool doProofPeriodMonitoring)
public ChainState(ILog log, ICodexContracts contracts, IChainStateChangeHandler changeHandler, DateTime startUtc)
{
this.log = new LogPrefixer(log, "(ChainState) ");
this.contracts = contracts;
handler = changeHandler;
this.doProofPeriodMonitoring = doProofPeriodMonitoring;
TotalSpan = new TimeRange(startUtc, startUtc);
PeriodMonitor = new PeriodMonitor(this.log, contracts);
}
public TimeRange TotalSpan { get; private set; }
public IChainStateRequest[] Requests => requests.ToArray();
public PeriodMonitor PeriodMonitor { get; }
public int Update()
{
@@ -92,18 +88,11 @@ namespace CodexContractsPlugin.ChainMonitor
{
var blockEvents = events.All.Where(e => e.Block.BlockNumber == b).ToArray();
ApplyEvents(b, blockEvents, eventUtc);
UpdatePeriodMonitor(b, eventUtc);
eventUtc += spanPerBlock;
}
}
private void UpdatePeriodMonitor(ulong blockNumber, DateTime eventUtc)
{
if (!doProofPeriodMonitoring) return;
PeriodMonitor.Update(blockNumber, eventUtc, Requests);
}
private void ApplyEvents(ulong blockNumber, IHasBlock[] blockEvents, DateTime eventsUtc)
{
foreach (var e in blockEvents)
@@ -176,13 +165,6 @@ namespace CodexContractsPlugin.ChainMonitor
handler.OnSlotReservationsFull(new RequestEvent(@event.Block, r), @event.SlotIndex);
}
private void ApplyEvent(ProofSubmittedEventDTO @event)
{
var id = Base58.Encode(@event.Id);
log.Log($"[{@event.Block.BlockNumber}] Proof submitted (id:{id})");
handler.OnProofSubmitted(@event.Block, id);
}
private void ApplyTimeImplicitEvents(ulong blockNumber, DateTime eventsUtc)
{
foreach (var r in requests)
@@ -1,6 +1,5 @@
using BlockchainUtils;
using GethPlugin;
using System.Numerics;
using Utils;
namespace CodexContractsPlugin.ChainMonitor
{
@@ -57,10 +56,5 @@ namespace CodexContractsPlugin.ChainMonitor
{
foreach (var handler in Handlers) handler.OnError(msg);
}
public void OnProofSubmitted(BlockTimeEntry block, string id)
{
foreach (var handler in Handlers) handler.OnProofSubmitted(block, id);
}
}
}
@@ -1,6 +1,6 @@
using CodexContractsPlugin.Marketplace;
using GethPlugin;
using Logging;
using Utils;
namespace CodexContractsPlugin.ChainMonitor
{
@@ -1,6 +1,5 @@
using BlockchainUtils;
using GethPlugin;
using System.Numerics;
using Utils;
namespace CodexContractsPlugin.ChainMonitor
{
@@ -41,9 +40,5 @@ namespace CodexContractsPlugin.ChainMonitor
public void OnError(string msg)
{
}
public void OnProofSubmitted(BlockTimeEntry block, string id)
{
}
}
}
@@ -1,99 +0,0 @@
using Logging;
using Utils;
namespace CodexContractsPlugin.ChainMonitor
{
public class PeriodMonitor
{
private readonly ILog log;
private readonly ICodexContracts contracts;
private readonly List<PeriodReport> reports = new List<PeriodReport>();
private ulong? currentPeriod = null;
public PeriodMonitor(ILog log, ICodexContracts contracts)
{
this.log = log;
this.contracts = contracts;
}
public void Update(ulong blockNumber, DateTime eventUtc, IChainStateRequest[] requests)
{
var period = contracts.GetPeriodNumber(eventUtc);
if (!currentPeriod.HasValue)
{
currentPeriod = period;
return;
}
if (period == currentPeriod.Value) return;
CreateReportForPeriod(blockNumber - 1, currentPeriod.Value, requests);
currentPeriod = period;
}
public PeriodReport[] GetAndClearReports()
{
var result = reports.ToArray();
reports.Clear();
return result;
}
private void CreateReportForPeriod(ulong lastBlockInPeriod, ulong periodNumber, IChainStateRequest[] requests)
{
log.Log("Creating report for period " + periodNumber);
ulong total = 0;
ulong required = 0;
var missed = new List<PeriodProofMissed>();
foreach (var request in requests)
{
for (ulong slotIndex = 0; slotIndex < request.Request.Ask.Slots; slotIndex++)
{
var state = contracts.GetProofState(request.Request, slotIndex, lastBlockInPeriod, periodNumber);
total++;
if (state.Required)
{
required++;
if (state.Missing)
{
var idx = Convert.ToInt32(slotIndex);
var host = request.Hosts.GetHost(idx);
missed.Add(new PeriodProofMissed(host, request, idx));
}
}
}
}
reports.Add(new PeriodReport(periodNumber, total, required, missed.ToArray()));
}
}
public class PeriodReport
{
public PeriodReport(ulong periodNumber, ulong totalNumSlots, ulong totalProofsRequired, PeriodProofMissed[] missedProofs)
{
PeriodNumber = periodNumber;
TotalNumSlots = totalNumSlots;
TotalProofsRequired = totalProofsRequired;
MissedProofs = missedProofs;
}
public ulong PeriodNumber { get; }
public ulong TotalNumSlots { get; }
public ulong TotalProofsRequired { get; }
public PeriodProofMissed[] MissedProofs { get; }
}
public class PeriodProofMissed
{
public PeriodProofMissed(EthAddress? host, IChainStateRequest request, int slotIndex)
{
Host = host;
Request = request;
SlotIndex = slotIndex;
}
public EthAddress? Host { get; }
public IChainStateRequest Request { get; }
public int SlotIndex { get; }
}
}
@@ -1,10 +1,7 @@
using BlockchainUtils;
using CodexContractsPlugin.Marketplace;
using CodexContractsPlugin.Marketplace;
using GethPlugin;
using Logging;
using Nethereum.ABI;
using Nethereum.ABI.FunctionEncoding.Attributes;
using Nethereum.Contracts;
using Nethereum.Hex.HexConvertors.Extensions;
using Nethereum.Util;
using NethereumWorkflow;
@@ -28,21 +25,7 @@ namespace CodexContractsPlugin
ICodexContractsEvents GetEvents(BlockInterval blockInterval);
EthAddress? GetSlotHost(Request storageRequest, decimal slotIndex);
RequestState GetRequestState(Request request);
ulong GetPeriodNumber(DateTime utc);
void WaitUntilNextPeriod();
ProofState GetProofState(Request storageRequest, decimal slotIndex, ulong blockNumber, ulong period);
}
public class ProofState
{
public ProofState(bool required, bool missing)
{
Required = required;
Missing = missing;
}
public bool Required { get; }
public bool Missing { get; }
}
[JsonConverter(typeof(StringEnumConverter))]
@@ -105,23 +88,19 @@ namespace CodexContractsPlugin
return new CodexContractsEvents(log, gethNode, Deployment, blockInterval);
}
public byte[] GetSlotId(Request request, decimal slotIndex)
public EthAddress? GetSlotHost(Request storageRequest, decimal slotIndex)
{
var encoder = new ABIEncode();
var encoded = encoder.GetABIEncoded(
new ABIValue("bytes32", request.RequestId),
new ABIValue("bytes32", storageRequest.RequestId),
new ABIValue("uint256", slotIndex.ToBig())
);
return Sha3Keccack.Current.CalculateHash(encoded);
}
var hashed = Sha3Keccack.Current.CalculateHash(encoded);
public EthAddress? GetSlotHost(Request storageRequest, decimal slotIndex)
{
var slotId = GetSlotId(storageRequest, slotIndex);
var func = new GetHostFunction
{
SlotId = slotId
SlotId = hashed
};
var address = gethNode.Call<GetHostFunction, string>(Deployment.MarketplaceAddress, func);
if (string.IsNullOrEmpty(address)) return null;
@@ -137,15 +116,6 @@ namespace CodexContractsPlugin
return gethNode.Call<RequestStateFunction, RequestState>(Deployment.MarketplaceAddress, func);
}
public ulong GetPeriodNumber(DateTime utc)
{
DateTimeOffset utco = DateTime.SpecifyKind(utc, DateTimeKind.Utc);
var now = utco.ToUnixTimeSeconds();
var periodSeconds = (int)Deployment.Config.Proofs.Period;
var result = now / periodSeconds;
return Convert.ToUInt64(result);
}
public void WaitUntilNextPeriod()
{
log.Log("Waiting until next proof period...");
@@ -155,52 +125,6 @@ namespace CodexContractsPlugin
Thread.Sleep(TimeSpan.FromSeconds(secondsLeft + 1));
}
public ProofState GetProofState(Request storageRequest, decimal slotIndex, ulong blockNumber, ulong period)
{
var slotId = GetSlotId(storageRequest, slotIndex);
var required = IsProofRequired(slotId, blockNumber);
if (!required) return new ProofState(false, false);
var missing = IsProofMissing(slotId, blockNumber, period);
return new ProofState(required, missing);
}
private bool IsProofRequired(byte[] slotId, ulong blockNumber)
{
var func = new IsProofRequiredFunction
{
Id = slotId
};
var result = gethNode.Call<IsProofRequiredFunction, IsProofRequiredOutputDTO>(Deployment.MarketplaceAddress, func, blockNumber);
return result.ReturnValue1;
}
private bool IsProofMissing(byte[] slotId, ulong blockNumber, ulong period)
{
try
{
var funcB = new MarkProofAsMissingFunction
{
SlotId = slotId,
Period = period
};
gethNode.Call<MarkProofAsMissingFunction>(Deployment.MarketplaceAddress, funcB, blockNumber);
}
catch (AggregateException exc)
{
if (exc.InnerExceptions.Count == 1)
{
if (exc.InnerExceptions[0].GetType() == typeof(SmartContractCustomErrorRevertException))
{
return false;
}
}
throw;
}
return true;
}
private ContractInteractions StartInteraction()
{
return new ContractInteractions(log, gethNode);
@@ -1,6 +1,7 @@
using GethPlugin;
using KubernetesWorkflow;
using KubernetesWorkflow.Recipe;
using Logging;
namespace CodexContractsPlugin
{
@@ -1,9 +1,9 @@
using BlockchainUtils;
using CodexContractsPlugin.Marketplace;
using CodexContractsPlugin.Marketplace;
using GethPlugin;
using Logging;
using Nethereum.Contracts;
using Nethereum.Hex.HexTypes;
using NethereumWorkflow.BlockUtils;
using Utils;
namespace CodexContractsPlugin
@@ -17,8 +17,7 @@ namespace CodexContractsPlugin
RequestFailedEventDTO[] GetRequestFailedEvents();
SlotFilledEventDTO[] GetSlotFilledEvents();
SlotFreedEventDTO[] GetSlotFreedEvents();
SlotReservationsFullEventDTO[] GetSlotReservationsFullEvents();
ProofSubmittedEventDTO[] GetProofSubmittedEvents();
SlotReservationsFullEventDTO[] GetSlotReservationsFull();
}
public class CodexContractsEvents : ICodexContractsEvents
@@ -87,18 +86,12 @@ namespace CodexContractsPlugin
return events.Select(SetBlockOnEvent).ToArray();
}
public SlotReservationsFullEventDTO[] GetSlotReservationsFullEvents()
public SlotReservationsFullEventDTO[] GetSlotReservationsFull()
{
var events = gethNode.GetEvents<SlotReservationsFullEventDTO>(deployment.MarketplaceAddress, BlockInterval);
return events.Select(SetBlockOnEvent).ToArray();
}
public ProofSubmittedEventDTO[] GetProofSubmittedEvents()
{
var events = gethNode.GetEvents<ProofSubmittedEventDTO>(deployment.MarketplaceAddress, BlockInterval);
return events.Select(SetBlockOnEvent).ToArray();
}
private T SetBlockOnEvent<T>(EventLog<T> e) where T : IHasBlock
{
var result = e.Event;
@@ -1,12 +1,11 @@
using BlockchainUtils;
using CodexContractsPlugin.Marketplace;
using CodexContractsPlugin.Marketplace;
using GethPlugin;
using Logging;
using Nethereum.ABI.FunctionEncoding.Attributes;
using Nethereum.Contracts;
using Nethereum.Hex.HexConvertors.Extensions;
using NethereumWorkflow;
using System.Numerics;
using Utils;
namespace CodexContractsPlugin
{
@@ -1,7 +1,7 @@
#pragma warning disable CS8618 // Non-nullable field must contain a non-null value when exiting constructor. Consider declaring as nullable.
using BlockchainUtils;
using GethPlugin;
using NethereumWorkflow.BlockUtils;
using Newtonsoft.Json;
using Utils;
namespace CodexContractsPlugin.Marketplace
{
@@ -69,11 +69,5 @@ namespace CodexContractsPlugin.Marketplace
[JsonIgnore]
public BlockTimeEntry Block { get; set; }
}
public partial class ProofSubmittedEventDTO : IHasBlock
{
[JsonIgnore]
public BlockTimeEntry Block { get; set; }
}
}
#pragma warning restore CS8618 // Non-nullable field must contain a non-null value when exiting constructor. Consider declaring as nullable.
File diff suppressed because one or more lines are too long
@@ -1,6 +1,6 @@
using System.Numerics;
namespace Utils
namespace CodexContractsPlugin
{
public class TestToken : IComparable<TestToken>
{
@@ -81,7 +81,7 @@ namespace Utils
{
public static TestToken TstWei(this int i)
{
return new TestToken(new BigInteger(i));
return TstWei(Convert.ToDecimal(i));
}
public static TestToken TstWei(this decimal i)
@@ -94,11 +94,6 @@ namespace Utils
return new TestToken(i);
}
public static TestToken TstWei(this string s)
{
return new TestToken(BigInteger.Parse(s));
}
public static TestToken Tst(this int i)
{
return Tst(Convert.ToDecimal(i));
@@ -113,10 +108,5 @@ namespace Utils
{
return new TestToken(i * TestToken.WeiFactor);
}
public static TestToken Tst(this string s)
{
return new TestToken(BigInteger.Parse(s) * TestToken.WeiFactor);
}
}
}
+4 -4
View File
@@ -10,7 +10,7 @@ namespace CodexPlugin
public class ApiChecker
{
// <INSERT-OPENAPI-YAML-HASH>
private const string OpenApiYamlHash = "00-7D-C3-0B-D2-23-D6-6C-CA-C2-43-D0-9B-B4-63-FC-4F-FE-23-9F-B8-82-5F-3B-3F-6B-4F-1F-11-E9-48-16";
private const string OpenApiYamlHash = "34-B5-DA-26-40-76-B8-D8-8E-7D-9C-17-85-C6-B0-63-55-8D-C6-01-0B-96-BB-7C-BD-53-E5-32-07-ED-29-92";
private const string OpenApiFilePath = "/codex/openapi.yaml";
private const string DisableEnvironmentVariable = "CODEXPLUGIN_DISABLE_APICHECK";
@@ -82,12 +82,12 @@ namespace CodexPlugin
private void OverwriteOpenApiYaml(string containerApi)
{
Log("API compatibility check failed. Updating CodexPlugin...");
var openApiFilePath = Path.Combine(PluginPathUtils.ProjectPluginsDir, "CodexClient", "openapi.yaml");
if (!File.Exists(openApiFilePath)) throw new Exception("Unable to locate CodexClient/openapi.yaml. Expected: " + openApiFilePath);
var openApiFilePath = Path.Combine(PluginPathUtils.ProjectPluginsDir, "CodexPlugin", "openapi.yaml");
if (!File.Exists(openApiFilePath)) throw new Exception("Unable to locate CodexPlugin/openapi.yaml. Expected: " + openApiFilePath);
File.Delete(openApiFilePath);
File.WriteAllText(openApiFilePath, containerApi);
Log("CodexClient/openapi.yaml has been updated.");
Log("CodexPlugin/openapi.yaml has been updated.");
}
private string Hash(string file)
@@ -1,156 +0,0 @@
using CodexClient;
using Core;
using Utils;
using System.Diagnostics;
namespace CodexPlugin
{
public class BinaryCodexStarter : ICodexStarter
{
private readonly IPluginTools pluginTools;
private readonly ProcessControlMap processControlMap;
private readonly static NumberSource numberSource = new NumberSource(1);
private readonly static FreePortFinder freePortFinder = new FreePortFinder();
private readonly static object _lock = new object();
private readonly static string dataParentDir = "codex_disttest_datadirs";
private readonly static CodexExePath codexExePath = new CodexExePath();
static BinaryCodexStarter()
{
StopAllCodexProcesses();
DeleteParentDataDir();
}
public BinaryCodexStarter(IPluginTools pluginTools, ProcessControlMap processControlMap)
{
this.pluginTools = pluginTools;
this.processControlMap = processControlMap;
}
public ICodexInstance[] BringOnline(CodexSetup codexSetup)
{
lock (_lock)
{
LogSeparator();
Log($"Starting {codexSetup.Describe()}...");
return StartCodexBinaries(codexSetup, codexSetup.NumberOfNodes);
}
}
public void Decommission()
{
lock (_lock)
{
processControlMap.StopAll();
}
}
private ICodexInstance[] StartCodexBinaries(CodexStartupConfig startupConfig, int numberOfNodes)
{
var result = new List<ICodexInstance>();
for (var i = 0; i < numberOfNodes; i++)
{
result.Add(StartBinary(startupConfig));
}
return result.ToArray();
}
private ICodexInstance StartBinary(CodexStartupConfig config)
{
var name = GetName(config);
var dataDir = Path.Combine(dataParentDir, $"datadir_{numberSource.GetNextNumber()}");
var pconfig = new CodexProcessConfig(name, freePortFinder, dataDir);
Log(pconfig);
var factory = new CodexProcessRecipe(pconfig, codexExePath);
var recipe = factory.Initialize(config);
var startInfo = new ProcessStartInfo(
fileName: recipe.Cmd,
arguments: recipe.Args
);
//startInfo.UseShellExecute = true;
startInfo.RedirectStandardOutput = true;
startInfo.RedirectStandardError = true;
var process = Process.Start(startInfo);
if (process == null || process.HasExited)
{
throw new Exception("Failed to start");
}
var local = "localhost";
var instance = new CodexInstance(
name: name,
imageName: "binary",
startUtc: DateTime.UtcNow,
discoveryEndpoint: new Address("Disc", pconfig.LocalIpAddrs.ToString(), pconfig.DiscPort),
apiEndpoint: new Address("Api", "http://" + local, pconfig.ApiPort),
listenEndpoint: new Address("Listen", local, pconfig.ListenPort),
ethAccount: null,
metricsEndpoint: null
);
var pc = new BinaryProcessControl(pluginTools.GetLog(), process, pconfig);
processControlMap.Add(instance, pc);
return instance;
}
private string GetName(CodexStartupConfig config)
{
if (!string.IsNullOrEmpty(config.NameOverride))
{
return config.NameOverride + "_" + numberSource.GetNextNumber();
}
return "codex_" + numberSource.GetNextNumber();
}
private void LogSeparator()
{
Log("----------------------------------------------------------------------------");
}
private void Log(CodexProcessConfig pconfig)
{
Log(
"NodeConfig:Name=" + pconfig.Name +
"ApiPort=" + pconfig.ApiPort +
"DiscPort=" + pconfig.DiscPort +
"ListenPort=" + pconfig.ListenPort +
"DataDir=" + pconfig.DataDir
);
}
private void Log(string message)
{
pluginTools.GetLog().Log(message);
}
private static void DeleteParentDataDir()
{
if (Directory.Exists(dataParentDir))
{
Directory.Delete(dataParentDir, true);
}
}
private static void StopAllCodexProcesses()
{
var processes = Process.GetProcesses();
var codexes = processes.Where(p =>
p.ProcessName.ToLowerInvariant() == "codex" &&
p.MainModule != null &&
p.MainModule.FileName == codexExePath.Get()
).ToArray();
foreach (var c in codexes)
{
c.Kill();
c.WaitForExit();
}
}
}
}
@@ -1,93 +0,0 @@
using System.Diagnostics;
using CodexClient;
using Logging;
namespace CodexPlugin
{
public class BinaryProcessControl : IProcessControl
{
private readonly LogFile logFile;
private readonly Process process;
private readonly CodexProcessConfig config;
private List<string> logBuffer = new List<string>();
private readonly object bufferLock = new object();
private readonly List<Task> streamTasks = new List<Task>();
private bool running;
public BinaryProcessControl(ILog log, Process process, CodexProcessConfig config)
{
logFile = log.CreateSubfile(config.Name);
running = true;
this.process = process;
this.config = config;
streamTasks.Add(Task.Run(() => ReadProcessStream(process.StandardOutput)));
streamTasks.Add(Task.Run(() => ReadProcessStream(process.StandardError)));
streamTasks.Add(Task.Run(() => WriteLog()));
}
private void ReadProcessStream(StreamReader reader)
{
while (running)
{
var line = reader.ReadLine();
if (!string.IsNullOrEmpty(line))
{
lock (bufferLock)
{
logBuffer.Add(line);
}
}
}
}
private void WriteLog()
{
while (running || logBuffer.Count > 0)
{
if (logBuffer.Count > 0)
{
List<string> lines = null!;
lock (bufferLock)
{
lines = logBuffer;
logBuffer = new List<string>();
}
logFile.WriteRawMany(lines);
}
else Thread.Sleep(100);
}
}
public void DeleteDataDirFolder()
{
if (!Directory.Exists(config.DataDir)) throw new Exception("datadir not found");
Directory.Delete(config.DataDir, true);
}
public IDownloadedLog DownloadLog(LogFile file)
{
return new DownloadedLog(logFile, config.Name);
}
public bool HasCrashed()
{
return process.HasExited;
}
public void Stop(bool waitTillStopped)
{
running = false;
process.Kill();
if (waitTillStopped)
{
process.WaitForExit();
foreach (var t in streamTasks) t.Wait();
}
DeleteDataDirFolder();
}
}
}
@@ -1,72 +1,37 @@
using CodexOpenApi;
using Core;
using KubernetesWorkflow;
using KubernetesWorkflow.Types;
using Logging;
using Newtonsoft.Json;
using Utils;
using WebUtils;
namespace CodexClient
namespace CodexPlugin
{
public class CodexAccess
{
private readonly ILog log;
private readonly IHttpFactory httpFactory;
private readonly IProcessControl processControl;
private ICodexInstance instance;
private readonly IPluginTools tools;
private readonly Mapper mapper = new Mapper();
public CodexAccess(ILog log, IHttpFactory httpFactory, IProcessControl processControl, ICodexInstance instance)
public CodexAccess(IPluginTools tools, RunningPod container, CrashWatcher crashWatcher)
{
this.log = log;
this.httpFactory = httpFactory;
this.processControl = processControl;
this.instance = instance;
this.tools = tools;
log = tools.GetLog();
Container = container;
CrashWatcher = crashWatcher;
CrashWatcher.Start();
}
public void Stop(bool waitTillStopped)
{
processControl.Stop(waitTillStopped);
// Prevents accidental use after stop:
instance = null!;
}
public IDownloadedLog DownloadLog(string additionalName = "")
{
var file = log.CreateSubfile(GetName() + additionalName);
Log($"Downloading logs to '{file.Filename}'");
return processControl.DownloadLog(file);
}
public string GetImageName()
{
return instance.ImageName;
}
public DateTime GetStartUtc()
{
return instance.StartUtc;
}
public RunningPod Container { get; }
public CrashWatcher CrashWatcher { get; }
public DebugInfo GetDebugInfo()
{
return mapper.Map(OnCodex(api => api.GetDebugInfoAsync()));
}
public void SetLogLevel(string logLevel)
{
try
{
OnCodex(async api =>
{
await api.SetDebugLogLevelAsync(logLevel);
return string.Empty;
});
}
catch (Exception exc)
{
log.Error("Failed to set log level: " + exc);
}
}
public string GetSpr()
{
return CrashCheck(() =>
@@ -114,14 +79,19 @@ namespace CodexClient
});
}
public string UploadFile(UploadInput uploadInput)
public string UploadFile(UploadInput uploadInput, Action<Failure> onFailure)
{
return OnCodex(api => api.UploadAsync(uploadInput.ContentType, uploadInput.ContentDisposition, uploadInput.FileStream));
return OnCodex(
api => api.UploadAsync(uploadInput.ContentType, uploadInput.ContentDisposition, uploadInput.FileStream),
CreateRetryConfig(nameof(UploadFile), onFailure));
}
public Stream DownloadFile(string contentId)
public Stream DownloadFile(string contentId, Action<Failure> onFailure)
{
var fileResponse = OnCodex(api => api.DownloadNetworkStreamAsync(contentId));
var fileResponse = OnCodex(
api => api.DownloadNetworkStreamAsync(contentId),
CreateRetryConfig(nameof(DownloadFile), onFailure));
if (fileResponse.StatusCode != 200) throw new Exception("Download failed with StatusCode: " + fileResponse.StatusCode);
return fileResponse.Stream;
}
@@ -168,25 +138,17 @@ namespace CodexClient
return mapper.Map(space);
}
public StoragePurchase? GetPurchaseStatus(string purchaseId)
public StoragePurchase GetPurchaseStatus(string purchaseId)
{
return CrashCheck(() =>
{
var endpoint = GetEndpoint();
try
return Time.Retry(() =>
{
return Time.Retry(() =>
{
var str = endpoint.HttpGetString($"storage/purchases/{purchaseId}");
if (string.IsNullOrEmpty(str)) throw new Exception("Empty response.");
return JsonConvert.DeserializeObject<StoragePurchase>(str)!;
}, nameof(GetPurchaseStatus));
}
catch (Exception exc)
{
log.Error($"Failed to fetch purchase information for id: '{purchaseId}'. Exception: {exc.Message}");
return null;
}
var str = endpoint.HttpGetString($"storage/purchases/{purchaseId}");
if (string.IsNullOrEmpty(str)) throw new Exception("Empty response.");
return JsonConvert.DeserializeObject<StoragePurchase>(str)!;
}, nameof(GetPurchaseStatus));
});
// TODO: current getpurchase api does not line up with its openapi spec.
@@ -195,60 +157,47 @@ namespace CodexClient
public string GetName()
{
return instance.Name;
return Container.Name;
}
public Address GetDiscoveryEndpoint()
public PodInfo GetPodInfo()
{
return instance.DiscoveryEndpoint;
var workflow = tools.CreateWorkflow();
return workflow.GetPodInfo(Container);
}
public Address GetApiEndpoint()
public void DeleteRepoFolder()
{
return instance.ApiEndpoint;
try
{
var containerNumber = Container.Containers.First().Recipe.Number;
var dataDir = $"datadir{containerNumber}";
var workflow = tools.CreateWorkflow();
workflow.ExecuteCommand(Container.Containers.First(), "rm", "-Rfv", $"/codex/{dataDir}/repo");
Log("Deleted repo folder.");
}
catch (Exception e)
{
Log("Unable to delete repo folder: " + e);
}
}
public Address GetListenEndpoint()
private T OnCodex<T>(Func<CodexApi, Task<T>> action)
{
return instance.ListenEndpoint;
}
public bool HasCrashed()
{
return processControl.HasCrashed();
}
public Address? GetMetricsEndpoint()
{
return instance.MetricsEndpoint;
}
public EthAccount? GetEthAccount()
{
return instance.EthAccount;
}
public void DeleteDataDirFolder()
{
processControl.DeleteDataDirFolder();
}
private T OnCodex<T>(Func<CodexApiClient, Task<T>> action)
{
var result = httpFactory.CreateHttp(GetHttpId(), h => CheckContainerCrashed()).OnClient(client => CallCodex(client, action));
var result = tools.CreateHttp(GetHttpId(), CheckContainerCrashed).OnClient(client => CallCodex(client, action));
return result;
}
private T OnCodex<T>(Func<CodexApiClient, Task<T>> action, Retry retry)
private T OnCodex<T>(Func<CodexApi, Task<T>> action, Retry retry)
{
var result = httpFactory.CreateHttp(GetHttpId(), h => CheckContainerCrashed()).OnClient(client => CallCodex(client, action), retry);
var result = tools.CreateHttp(GetHttpId(), CheckContainerCrashed).OnClient(client => CallCodex(client, action), retry);
return result;
}
private T CallCodex<T>(HttpClient client, Func<CodexApiClient, Task<T>> action)
private T CallCodex<T>(HttpClient client, Func<CodexApi, Task<T>> action)
{
var address = GetAddress();
var api = new CodexApiClient(client);
var api = new CodexApi(client);
api.BaseUrl = $"{address.Host}:{address.Port}/api/codex/v1";
return CrashCheck(() => Time.Wait(action(api)));
}
@@ -261,20 +210,20 @@ namespace CodexClient
}
finally
{
CheckContainerCrashed();
CrashWatcher.HasContainerCrashed();
}
}
private IEndpoint GetEndpoint()
{
return httpFactory
.CreateHttp(GetHttpId(), h => CheckContainerCrashed())
.CreateEndpoint(GetAddress(), "/api/codex/v1/", GetName());
return tools
.CreateHttp(GetHttpId(), CheckContainerCrashed)
.CreateEndpoint(GetAddress(), "/api/codex/v1/", Container.Name);
}
private Address GetAddress()
{
return instance.ApiEndpoint;
return Container.Containers.Single().GetAddress(CodexContainerRecipe.ApiPortTag);
}
private string GetHttpId()
@@ -282,9 +231,52 @@ namespace CodexClient
return GetAddress().ToString();
}
private void CheckContainerCrashed()
private void CheckContainerCrashed(HttpClient client)
{
if (processControl.HasCrashed()) throw new Exception($"Container {GetName()} has crashed.");
if (CrashWatcher.HasContainerCrashed()) throw new Exception($"Container {GetName()} has crashed.");
}
private Retry CreateRetryConfig(string description, Action<Failure> onFailure)
{
var timeSet = tools.TimeSet;
return new Retry(description, timeSet.HttpRetryTimeout(), timeSet.HttpCallRetryDelay(), failure =>
{
onFailure(failure);
Investigate(failure, timeSet);
});
}
private void Investigate(Failure failure, ITimeSet timeSet)
{
Log($"Retry {failure.TryNumber} took {Time.FormatDuration(failure.Duration)} and failed with '{failure.Exception}'. " +
$"(HTTP timeout = {Time.FormatDuration(timeSet.HttpCallTimeout())}) " +
$"Checking if node responds to debug/info...");
try
{
var debugInfo = GetDebugInfo();
if (string.IsNullOrEmpty(debugInfo.Spr))
{
Log("Did not get value debug/info response.");
Throw(failure);
}
else
{
Log("Got valid response from debug/info.");
}
}
catch (Exception ex)
{
Log("Got exception from debug/info call: " + ex);
Throw(failure);
}
if (failure.Duration < timeSet.HttpCallTimeout())
{
Log("Retry failed within HTTP timeout duration.");
Throw(failure);
}
}
private void Throw(Failure failure)
@@ -294,7 +286,7 @@ namespace CodexClient
private void Log(string msg)
{
log.Log($"({GetName()}) {msg}");
log.Log($"{GetName()} {msg}");
}
}
@@ -1,70 +0,0 @@
using CodexClient;
using Core;
using KubernetesWorkflow;
using KubernetesWorkflow.Types;
using Logging;
namespace CodexPlugin
{
public class CodexContainerProcessControl : IProcessControl
{
private readonly IPluginTools tools;
private readonly RunningPod pod;
private readonly Action onStop;
private readonly ContainerCrashWatcher crashWatcher;
public CodexContainerProcessControl(IPluginTools tools, RunningPod pod, Action onStop)
{
this.tools = tools;
this.pod = pod;
this.onStop = onStop;
crashWatcher = tools.CreateWorkflow().CreateCrashWatcher(pod.Containers.Single());
crashWatcher.Start();
}
public void Stop(bool waitTillStopped)
{
Log($"Stopping node...");
crashWatcher.Stop();
var workflow = tools.CreateWorkflow();
workflow.Stop(pod, waitTillStopped);
onStop();
Log("Stopped.");
}
public IDownloadedLog DownloadLog(LogFile file)
{
var workflow = tools.CreateWorkflow();
return workflow.DownloadContainerLog(pod.Containers.Single());
}
public void DeleteDataDirFolder()
{
var container = pod.Containers.Single();
try
{
var dataDirVar = container.Recipe.EnvVars.Single(e => e.Name == "CODEX_DATA_DIR");
var dataDir = dataDirVar.Value;
var workflow = tools.CreateWorkflow();
workflow.ExecuteCommand(container, "rm", "-Rfv", $"/codex/{dataDir}/repo");
Log("Deleted repo folder.");
}
catch (Exception e)
{
Log("Unable to delete repo folder: " + e);
}
}
public bool HasCrashed()
{
return crashWatcher.HasCrashed();
}
private void Log(string message)
{
tools.GetLog().Log(message);
}
}
}
+15 -4
View File
@@ -1,5 +1,4 @@
using CodexClient;
using CodexContractsPlugin;
using CodexContractsPlugin;
using GethPlugin;
using KubernetesWorkflow.Types;
@@ -10,7 +9,7 @@ namespace CodexPlugin
public CodexDeployment(CodexInstance[] codexInstances, GethDeployment gethDeployment,
CodexContractsDeployment codexContractsDeployment, RunningPod? prometheusContainer,
RunningPod? discordBotContainer, DeploymentMetadata metadata,
string id)
String id)
{
Id = id;
CodexInstances = codexInstances;
@@ -21,7 +20,7 @@ namespace CodexPlugin
Metadata = metadata;
}
public string Id { get; }
public String Id { get; }
public CodexInstance[] CodexInstances { get; }
public GethDeployment GethDeployment { get; }
public CodexContractsDeployment CodexContractsDeployment { get; }
@@ -30,6 +29,18 @@ namespace CodexPlugin
public DeploymentMetadata Metadata { get; }
}
public class CodexInstance
{
public CodexInstance(RunningPod pod, DebugInfo info)
{
Pod = pod;
Info = info;
}
public RunningPod Pod { get; }
public DebugInfo Info { get; }
}
public class DeploymentMetadata
{
public DeploymentMetadata(string name, DateTime startUtc, DateTime finishedUtc, string kubeNamespace,
@@ -1,29 +0,0 @@
namespace CodexPlugin
{
public class CodexExePath
{
private readonly string[] paths = [
Path.Combine("d:", "Dev", "nim-codex", "build", "codex.exe"),
Path.Combine("c:", "Projects", "nim-codex", "build", "codex.exe")
];
private string selectedPath = string.Empty;
public CodexExePath()
{
foreach (var p in paths)
{
if (File.Exists(p))
{
selectedPath = p;
return;
}
}
}
public string Get()
{
return selectedPath;
}
}
}
@@ -1,46 +0,0 @@
using CodexClient;
using KubernetesWorkflow.Types;
using Utils;
namespace CodexPlugin
{
public static class CodexInstanceContainerExtension
{
public static ICodexInstance CreateFromPod(RunningPod pod)
{
var container = pod.Containers.Single();
return new CodexInstance(
name: container.Name,
imageName: container.Recipe.Image,
startUtc: container.Recipe.RecipeCreatedUtc,
discoveryEndpoint: SetClusterInternalIpAddress(pod, container.GetInternalAddress(CodexContainerRecipe.DiscoveryPortTag)),
apiEndpoint: container.GetAddress(CodexContainerRecipe.ApiPortTag),
listenEndpoint: container.GetInternalAddress(CodexContainerRecipe.ListenPortTag),
ethAccount: container.Recipe.Additionals.Get<EthAccount>(),
metricsEndpoint: GetMetricsEndpoint(container)
);
}
private static Address SetClusterInternalIpAddress(RunningPod pod, Address address)
{
return new Address(
logName: address.LogName,
host: pod.PodInfo.Ip,
port: address.Port
);
}
private static Address? GetMetricsEndpoint(RunningContainer container)
{
try
{
return container.GetInternalAddress(CodexContainerRecipe.MetricsPortTag);
}
catch
{
return null;
}
}
}
}
@@ -1,4 +1,4 @@
namespace CodexClient
namespace CodexPlugin
{
public enum CodexLogLevel
{
@@ -1,6 +1,6 @@
using System.Globalization;
namespace CodexClient
namespace CodexPlugin
{
public class CodexLogLine
{
@@ -1,95 +1,98 @@
using CodexClient.Hooks;
using CodexPlugin.Hooks;
using Core;
using FileUtils;
using GethPlugin;
using KubernetesWorkflow;
using KubernetesWorkflow.Types;
using Logging;
using MetricsPlugin;
using Utils;
namespace CodexClient
namespace CodexPlugin
{
public partial interface ICodexNode : IHasEthAddress, IHasMetricsScrapeTarget
public interface ICodexNode : IHasContainer, IHasMetricsScrapeTarget, IHasEthAddress
{
string GetName();
string GetImageName();
string GetPeerId();
DebugInfo GetDebugInfo(bool log = false);
void SetLogLevel(string logLevel);
string GetSpr();
DebugPeer GetDebugPeer(string peerId);
ContentId UploadFile(TrackedFile file);
ContentId UploadFile(TrackedFile file, string contentType, string contentDisposition);
ContentId UploadFile(TrackedFile file, Action<Failure> onFailure);
ContentId UploadFile(TrackedFile file, string contentType, string contentDisposition, Action<Failure> onFailure);
TrackedFile? DownloadContent(ContentId contentId, string fileLabel = "");
TrackedFile? DownloadContent(ContentId contentId, Action<Failure> onFailure, string fileLabel = "");
LocalDataset DownloadStreamless(ContentId cid);
/// <summary>
/// TODO: This will monitor the quota-used of the node until 'size' bytes are added. That's a very bad way
/// to track the streamless download progress. Replace it once we have a good API for this.
/// </summary>
LocalDataset DownloadStreamlessWait(ContentId cid, ByteSize size);
LocalDataset DownloadManifestOnly(ContentId cid);
LocalDatasetList LocalFiles();
CodexSpace Space();
void ConnectToPeer(ICodexNode node);
DebugInfoVersion Version { get; }
IMarketplaceAccess Marketplace { get; }
CrashWatcher CrashWatcher { get; }
PodInfo GetPodInfo();
ITransferSpeeds TransferSpeeds { get; }
EthAccount EthAccount { get; }
StoragePurchase? GetPurchaseStatus(string purchaseId);
Address GetDiscoveryEndpoint();
Address GetApiEndpoint();
Address GetListenEndpoint();
/// <summary>
/// Warning! The node is not usable after this.
/// TODO: Replace with delete-blocks debug call once available in Codex.
/// </summary>
void DeleteDataDirFolder();
void DeleteRepoFolder();
void Stop(bool waitTillStopped);
IDownloadedLog DownloadLog(string additionalName = "");
bool HasCrashed();
}
public class CodexNode : ICodexNode
{
private const string UploadFailedMessage = "Unable to store block";
private readonly ILog log;
private readonly IPluginTools tools;
private readonly ICodexNodeHooks hooks;
private readonly EthAccount? ethAccount;
private readonly TransferSpeeds transferSpeeds;
private string peerId = string.Empty;
private string nodeId = string.Empty;
private readonly CodexAccess codexAccess;
private readonly IFileManager fileManager;
public CodexNode(ILog log, CodexAccess codexAccess, IFileManager fileManager, IMarketplaceAccess marketplaceAccess, ICodexNodeHooks hooks)
public CodexNode(IPluginTools tools, CodexAccess codexAccess, CodexNodeGroup group, IMarketplaceAccess marketplaceAccess, ICodexNodeHooks hooks, EthAccount? ethAccount)
{
this.codexAccess = codexAccess;
this.fileManager = fileManager;
this.tools = tools;
this.ethAccount = ethAccount;
CodexAccess = codexAccess;
Group = group;
Marketplace = marketplaceAccess;
this.hooks = hooks;
Version = new DebugInfoVersion();
transferSpeeds = new TransferSpeeds();
this.log = new LogPrefixer(log, $"{GetName()} ");
log = new LogPrefixer(tools.GetLog(), $"{GetName()} ");
}
public void Awake()
{
hooks.OnNodeStarting(codexAccess.GetStartUtc(), codexAccess.GetImageName(), codexAccess.GetEthAccount());
hooks.OnNodeStarting(Container.Recipe.RecipeCreatedUtc, Container.Recipe.Image, ethAccount);
}
public void Initialize()
{
InitializePeerNodeId();
InitializeLogReplacements();
hooks.OnNodeStarted(this, peerId, nodeId);
hooks.OnNodeStarted(peerId, nodeId);
}
public RunningPod Pod { get { return CodexAccess.Container; } }
public RunningContainer Container { get { return Pod.Containers.Single(); } }
public CodexAccess CodexAccess { get; }
public CrashWatcher CrashWatcher { get => CodexAccess.CrashWatcher; }
public CodexNodeGroup Group { get; }
public IMarketplaceAccess Marketplace { get; }
public DebugInfoVersion Version { get; private set; }
public ITransferSpeeds TransferSpeeds { get => transferSpeeds; }
public StoragePurchase? GetPurchaseStatus(string purchaseId)
public IMetricsScrapeTarget MetricsScrapeTarget
{
return codexAccess.GetPurchaseStatus(purchaseId);
get
{
return new MetricsScrapeTarget(CodexAccess.Container.Containers.First(), CodexContainerRecipe.MetricsPortTag);
}
}
public EthAddress EthAddress
@@ -97,7 +100,7 @@ namespace CodexClient
get
{
EnsureMarketplace();
return codexAccess.GetEthAccount()!.EthAddress;
return ethAccount!.EthAddress;
}
}
@@ -106,18 +109,13 @@ namespace CodexClient
get
{
EnsureMarketplace();
return codexAccess.GetEthAccount()!;
return ethAccount!;
}
}
public string GetName()
{
return codexAccess.GetName();
}
public string GetImageName()
{
return codexAccess.GetImageName();
return Container.Name;
}
public string GetPeerId()
@@ -127,7 +125,7 @@ namespace CodexClient
public DebugInfo GetDebugInfo(bool log = false)
{
var debugInfo = codexAccess.GetDebugInfo();
var debugInfo = CodexAccess.GetDebugInfo();
if (log)
{
var known = string.Join(",", debugInfo.Table.Nodes.Select(n => n.PeerId));
@@ -136,27 +134,27 @@ namespace CodexClient
return debugInfo;
}
public void SetLogLevel(string logLevel)
{
codexAccess.SetLogLevel(logLevel);
}
public string GetSpr()
{
return codexAccess.GetSpr();
return CodexAccess.GetSpr();
}
public DebugPeer GetDebugPeer(string peerId)
{
return codexAccess.GetDebugPeer(peerId);
return CodexAccess.GetDebugPeer(peerId);
}
public ContentId UploadFile(TrackedFile file)
{
return UploadFile(file, "application/octet-stream", $"attachment; filename=\"{Path.GetFileName(file.Filename)}\"");
return UploadFile(file, DoNothing);
}
public ContentId UploadFile(TrackedFile file, string contentType, string contentDisposition)
public ContentId UploadFile(TrackedFile file, Action<Failure> onFailure)
{
return UploadFile(file, "application/x-binary", $"attachment; filename=\"{file.Filename}\"", onFailure);
}
public ContentId UploadFile(TrackedFile file, string contentType, string contentDisposition, Action<Failure> onFailure)
{
using var fileStream = File.OpenRead(file.Filename);
var uniqueId = Guid.NewGuid().ToString();
@@ -168,7 +166,7 @@ namespace CodexClient
var logMessage = $"Uploading file {file.Describe()} with contentType: '{input.ContentType}' and disposition: '{input.ContentDisposition}'...";
var measurement = Stopwatch.Measure(log, logMessage, () =>
{
return codexAccess.UploadFile(input);
return CodexAccess.UploadFile(input, onFailure);
});
var response = measurement.Value;
@@ -186,12 +184,17 @@ namespace CodexClient
public TrackedFile? DownloadContent(ContentId contentId, string fileLabel = "")
{
var file = fileManager.CreateEmptyFile(fileLabel);
return DownloadContent(contentId, DoNothing, fileLabel);
}
public TrackedFile? DownloadContent(ContentId contentId, Action<Failure> onFailure, string fileLabel = "")
{
var file = tools.GetFileManager().CreateEmptyFile(fileLabel);
hooks.OnFileDownloading(contentId);
Log($"Downloading '{contentId}'...");
var logMessage = $"Downloaded '{contentId}' to '{file.Filename}'";
var measurement = Stopwatch.Measure(log, logMessage, () => DownloadToFile(contentId.Id, file));
var measurement = Stopwatch.Measure(log, logMessage, () => DownloadToFile(contentId.Id, file, onFailure));
var size = file.GetFilesize();
transferSpeeds.AddDownloadSample(size, measurement);
@@ -202,39 +205,22 @@ namespace CodexClient
public LocalDataset DownloadStreamless(ContentId cid)
{
Log($"Downloading streamless '{cid}' (no-wait)");
return codexAccess.DownloadStreamless(cid);
}
public LocalDataset DownloadStreamlessWait(ContentId cid, ByteSize size)
{
Log($"Downloading streamless '{cid}' (wait till finished)");
var sw = Stopwatch.Measure(log, nameof(DownloadStreamlessWait), () =>
{
var startSpace = Space();
var result = codexAccess.DownloadStreamless(cid);
WaitUntilQuotaUsedIncreased(startSpace, size);
return result;
});
return sw.Value;
return CodexAccess.DownloadStreamless(cid);
}
public LocalDataset DownloadManifestOnly(ContentId cid)
{
Log($"Downloading manifest-only '{cid}'");
return codexAccess.DownloadManifestOnly(cid);
return CodexAccess.DownloadManifestOnly(cid);
}
public LocalDatasetList LocalFiles()
{
return codexAccess.LocalFiles();
return CodexAccess.LocalFiles();
}
public CodexSpace Space()
{
return codexAccess.Space();
return CodexAccess.Space();
}
public void ConnectToPeer(ICodexNode node)
@@ -243,53 +229,47 @@ namespace CodexClient
Log($"Connecting to peer {peer.GetName()}...");
var peerInfo = node.GetDebugInfo();
codexAccess.ConnectToPeer(peerInfo.Id, GetPeerMultiAddresses(peer, peerInfo));
CodexAccess.ConnectToPeer(peerInfo.Id, GetPeerMultiAddresses(peer, peerInfo));
Log($"Successfully connected to peer {peer.GetName()}.");
}
public void DeleteDataDirFolder()
public PodInfo GetPodInfo()
{
codexAccess.DeleteDataDirFolder();
return CodexAccess.GetPodInfo();
}
public void DeleteRepoFolder()
{
CodexAccess.DeleteRepoFolder();
}
public void Stop(bool waitTillStopped)
{
Log("Stopping...");
hooks.OnNodeStopping();
codexAccess.Stop(waitTillStopped);
CrashWatcher.Stop();
Group.Stop(this, waitTillStopped);
}
public IDownloadedLog DownloadLog(string additionalName = "")
public void EnsureOnlineGetVersionResponse()
{
return codexAccess.DownloadLog(additionalName);
}
var debugInfo = Time.Retry(CodexAccess.GetDebugInfo, "ensure online");
peerId = debugInfo.Id;
nodeId = debugInfo.Table.LocalNode.NodeId;
var nodeName = CodexAccess.Container.Name;
public Address GetDiscoveryEndpoint()
{
return codexAccess.GetDiscoveryEndpoint();
}
if (!debugInfo.Version.IsValid())
{
throw new Exception($"Invalid version information received from Codex node {GetName()}: {debugInfo.Version}");
}
public Address GetApiEndpoint()
{
return codexAccess.GetApiEndpoint();
}
public Address GetListenEndpoint()
{
return codexAccess.GetListenEndpoint();
}
public Address GetMetricsScrapeTarget()
{
var address = codexAccess.GetMetricsEndpoint();
if (address == null) throw new Exception("Metrics ScrapeTarget accessed, but node was not started with EnableMetrics()");
return address;
}
public bool HasCrashed()
{
return codexAccess.HasCrashed();
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;
}
public override string ToString()
@@ -297,44 +277,22 @@ namespace CodexClient
return $"CodexNode:{GetName()}";
}
private void InitializePeerNodeId()
{
var debugInfo = Time.Retry(codexAccess.GetDebugInfo, "ensure online");
if (!debugInfo.Version.IsValid())
{
throw new Exception($"Invalid version information received from Codex node {GetName()}: {debugInfo.Version}");
}
peerId = debugInfo.Id;
nodeId = debugInfo.Table.LocalNode.NodeId;
Version = debugInfo.Version;
}
private void InitializeLogReplacements()
{
var nodeName = GetName();
log.AddStringReplace(peerId, nodeName);
log.AddStringReplace(CodexUtils.ToShortId(peerId), nodeName);
log.AddStringReplace(nodeId, nodeName);
log.AddStringReplace(CodexUtils.ToShortId(nodeId), nodeName);
}
private string[] GetPeerMultiAddresses(CodexNode peer, DebugInfo peerInfo)
{
var peerId = peer.GetDiscoveryEndpoint().Host
.Replace("http://", "")
.Replace("https://", "");
// The peer we want to connect is in a different pod.
// We must replace the default IP with the pod IP in the multiAddress.
var workflow = tools.CreateWorkflow();
var podInfo = workflow.GetPodInfo(peer.Pod);
return peerInfo.Addrs.Select(a => a
.Replace("0.0.0.0", peerId))
.Replace("0.0.0.0", podInfo.Ip))
.ToArray();
}
private void DownloadToFile(string contentId, TrackedFile file)
private void DownloadToFile(string contentId, TrackedFile file, Action<Failure> onFailure)
{
using var fileStream = File.OpenWrite(file.Filename);
var timeout = TimeSpan.FromMinutes(2.0); // todo: make this user-controllable.
var timeout = tools.TimeSet.HttpCallTimeout();
try
{
// Type of stream generated by openAPI client does not support timeouts.
@@ -342,7 +300,7 @@ namespace CodexClient
var cts = new CancellationTokenSource();
var downloadTask = Task.Run(() =>
{
using var downloadStream = codexAccess.DownloadFile(contentId);
using var downloadStream = CodexAccess.DownloadFile(contentId, onFailure);
downloadStream.CopyTo(fileStream);
}, cts.Token);
@@ -356,54 +314,25 @@ namespace CodexClient
cts.Cancel();
throw new TimeoutException($"Download of '{contentId}' timed out after {Time.FormatDuration(timeout)}");
}
catch (Exception ex)
catch
{
Log($"Failed to download file '{contentId}': {ex}");
Log($"Failed to download file '{contentId}'.");
throw;
}
}
public void WaitUntilQuotaUsedIncreased(CodexSpace startSpace, ByteSize expectedIncreaseOfQuotaUsed)
{
WaitUntilQuotaUsedIncreased(startSpace, expectedIncreaseOfQuotaUsed, TimeSpan.FromMinutes(2));
}
public void WaitUntilQuotaUsedIncreased(
CodexSpace startSpace,
ByteSize expectedIncreaseOfQuotaUsed,
TimeSpan maxTimeout)
{
Log($"Waiting until quotaUsed " +
$"(start: {startSpace.QuotaUsedBytes}) " +
$"increases by {expectedIncreaseOfQuotaUsed} " +
$"to reach {startSpace.QuotaUsedBytes + expectedIncreaseOfQuotaUsed.SizeInBytes}");
var retry = new Retry($"Checking local space for quotaUsed increase of {expectedIncreaseOfQuotaUsed}",
maxTimeout: maxTimeout,
sleepAfterFail: TimeSpan.FromSeconds(3),
onFail: f => { });
retry.Run(() =>
{
var space = Space();
var increase = space.QuotaUsedBytes - startSpace.QuotaUsedBytes;
if (increase < expectedIncreaseOfQuotaUsed.SizeInBytes)
throw new Exception($"Expected quota-used not reached. " +
$"Expected increase: {expectedIncreaseOfQuotaUsed.SizeInBytes} " +
$"Actual increase: {increase} " +
$"Actual used: {space.QuotaUsedBytes}");
});
}
private void EnsureMarketplace()
{
if (codexAccess.GetEthAccount() == null) throw new Exception("Marketplace is not enabled for this Codex node. Please start it with the option '.EnableMarketplace(...)' to enable it.");
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)
{
log.Log(msg);
}
private void DoNothing(Failure failure)
{
}
}
}
@@ -0,0 +1,53 @@
using CodexPlugin.Hooks;
using Core;
using GethPlugin;
using KubernetesWorkflow;
using KubernetesWorkflow.Types;
namespace CodexPlugin
{
public interface ICodexNodeFactory
{
CodexNode CreateOnlineCodexNode(CodexAccess access, CodexNodeGroup group);
CrashWatcher CreateCrashWatcher(RunningContainer c);
}
public class CodexNodeFactory : ICodexNodeFactory
{
private readonly IPluginTools tools;
private readonly CodexHooksFactory codexHooksFactory;
public CodexNodeFactory(IPluginTools tools, CodexHooksFactory codexHooksFactory)
{
this.tools = tools;
this.codexHooksFactory = codexHooksFactory;
}
public CodexNode CreateOnlineCodexNode(CodexAccess access, CodexNodeGroup group)
{
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, EthAccount? ethAccount, ICodexNodeHooks hooks)
{
if (ethAccount == null) return new MarketplaceUnavailable();
return new MarketplaceAccess(tools.GetLog(), codexAccess, hooks);
}
private EthAccount? GetEthAccount(CodexAccess access)
{
var ethAccount = access.Container.Containers.Single().Recipe.Additionals.Get<EthAccount>();
if (ethAccount == null) return null;
return ethAccount;
}
public CrashWatcher CreateCrashWatcher(RunningContainer c)
{
return tools.CreateWorkflow().CreateCrashWatcher(c);
}
}
}
+32 -17
View File
@@ -1,23 +1,25 @@
using CodexClient;
using Core;
using Core;
using KubernetesWorkflow.Types;
using MetricsPlugin;
using System.Collections;
using Utils;
namespace CodexPlugin
{
public interface ICodexNodeGroup : IEnumerable<ICodexNode>, IHasManyMetricScrapeTargets
{
void Stop(bool waitTillStopped);
void BringOffline(bool waitTillStopped);
ICodexNode this[int index] { get; }
}
public class CodexNodeGroup : ICodexNodeGroup
{
private readonly ICodexNode[] nodes;
private readonly CodexStarter starter;
public CodexNodeGroup(IPluginTools tools, ICodexNode[] nodes)
public CodexNodeGroup(CodexStarter starter, IPluginTools tools, RunningPod[] containers, ICodexNodeFactory codexNodeFactory)
{
this.nodes = nodes;
this.starter = starter;
Containers = containers;
Nodes = containers.Select(c => CreateOnlineCodexNode(c, tools, codexNodeFactory)).ToArray();
Version = new DebugInfoVersion();
}
@@ -29,23 +31,25 @@ namespace CodexPlugin
}
}
public void Stop(bool waitTillStopped)
public void BringOffline(bool waitTillStopped)
{
foreach (var node in Nodes) node.Stop(waitTillStopped);
starter.BringOffline(this, waitTillStopped);
// Clear everything. Prevent accidental use.
Nodes = Array.Empty<CodexNode>();
Containers = null!;
}
public void Stop(CodexNode node, bool waitTillStopped)
{
node.Stop(waitTillStopped);
starter.Stop(node.Pod, waitTillStopped);
Nodes = Nodes.Where(n => n != node).ToArray();
Containers = Containers.Where(c => c != node.Pod).ToArray();
}
public ICodexNode[] Nodes => nodes;
public RunningPod[] Containers { get; private set; }
public CodexNode[] Nodes { get; private set; }
public DebugInfoVersion Version { get; private set; }
public Address[] GetMetricsScrapeTargets()
{
return Nodes.Select(n => n.GetMetricsScrapeTarget()).ToArray();
}
public IMetricsScrapeTarget[] ScrapeTargets => Nodes.Select(n => n.MetricsScrapeTarget).ToArray();
public IEnumerator<ICodexNode> GetEnumerator()
{
@@ -59,11 +63,12 @@ namespace CodexPlugin
public string Describe()
{
return $"group:[{string.Join(",", Nodes.Select(n => n.GetName()))}]";
return $"group:[{Containers.Describe()}]";
}
public void EnsureOnline()
{
foreach (var node in Nodes) node.EnsureOnlineGetVersionResponse();
var versionResponses = Nodes.Select(n => n.Version);
var first = versionResponses.First();
@@ -74,6 +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);
var node = factory.CreateOnlineCodexNode(access, this);
node.Awake();
return node;
}
}
}
+13 -34
View File
@@ -1,68 +1,48 @@
using CodexClient;
using CodexClient.Hooks;
using CodexPlugin.Hooks;
using Core;
using KubernetesWorkflow.Types;
namespace CodexPlugin
{
public class CodexPlugin : IProjectPlugin, IHasLogPrefix, IHasMetadata
{
private const bool UseContainers = true;
private readonly ICodexStarter codexStarter;
private readonly CodexStarter codexStarter;
private readonly IPluginTools tools;
private readonly CodexLogLevel defaultLogLevel = CodexLogLevel.Trace;
private readonly CodexHooksFactory hooksFactory = new CodexHooksFactory();
private readonly ProcessControlMap processControlMap = new ProcessControlMap();
private readonly CodexWrapper codexWrapper;
public CodexPlugin(IPluginTools tools)
{
codexStarter = new CodexStarter(tools);
this.tools = tools;
codexStarter = CreateCodexStarter();
codexWrapper = new CodexWrapper(tools, processControlMap, hooksFactory);
}
private ICodexStarter CreateCodexStarter()
{
if (UseContainers)
{
Log("Using Containerized Codex instances");
return new ContainerCodexStarter(tools, processControlMap);
}
Log("Using Binary Codex instances");
return new BinaryCodexStarter(tools, processControlMap);
}
public string LogPrefix => "(Codex) ";
public void Announce()
{
Log($"Loaded with Codex ID: '{codexWrapper.GetCodexId()}' - Revision: {codexWrapper.GetCodexRevision()}");
Log($"Loaded with Codex ID: '{codexStarter.GetCodexId()}' - Revision: {codexStarter.GetCodexRevision()}");
}
public void AddMetadata(IAddMetadata metadata)
{
metadata.Add("codexid", codexWrapper.GetCodexId());
metadata.Add("codexrevision", codexWrapper.GetCodexRevision());
metadata.Add("codexid", codexStarter.GetCodexId());
metadata.Add("codexrevision", codexStarter.GetCodexRevision());
}
public void Decommission()
{
codexStarter.Decommission();
}
public ICodexInstance[] DeployCodexNodes(int numberOfNodes, Action<ICodexSetup> setup)
public RunningPod[] DeployCodexNodes(int numberOfNodes, Action<ICodexSetup> setup)
{
var codexSetup = GetSetup(numberOfNodes, setup);
return codexStarter.BringOnline(codexSetup);
}
public ICodexNodeGroup WrapCodexContainers(ICodexInstance[] instances)
public ICodexNodeGroup WrapCodexContainers(CoreInterface coreInterface, RunningPod[] containers)
{
instances = instances.Select(c => SerializeGate.Gate(c as CodexInstance)).ToArray();
return codexWrapper.WrapCodexInstances(instances);
containers = containers.Select(c => SerializeGate.Gate(c)).ToArray();
return codexStarter.WrapCodexContainers(coreInterface, containers);
}
public void WireUpMarketplace(ICodexNodeGroup result, Action<ICodexSetup> setup)
@@ -82,10 +62,9 @@ namespace CodexPlugin
}
}
public void AddCodexHooksProvider(ICodexHooksProvider hooksProvider)
public void SetCodexHooksProvider(ICodexHooksProvider hooksProvider)
{
if (hooksFactory.Providers.Contains(hooksProvider)) return;
hooksFactory.Providers.Add(hooksProvider);
codexStarter.HooksFactory.Provider = hooksProvider;
}
private CodexSetup GetSetup(int numberOfNodes, Action<ICodexSetup> setup)
@@ -1,4 +1,4 @@
<Project Sdk="Microsoft.NET.Sdk">
<Project Sdk="Microsoft.NET.Sdk">
<PropertyGroup>
<TargetFramework>net8.0</TargetFramework>
@@ -6,6 +6,14 @@
<Nullable>enable</Nullable>
</PropertyGroup>
<ItemGroup>
<None Remove="openapi.yaml" />
</ItemGroup>
<ItemGroup>
<OpenApiReference Include="openapi.yaml" CodeGenerator="NSwagCSharp" Namespace="CodexOpenApi" ClassName="CodexApi" />
</ItemGroup>
<ItemGroup>
<PackageReference Include="Microsoft.Extensions.ApiDescription.Client" Version="7.0.2">
<PrivateAssets>all</PrivateAssets>
@@ -22,7 +30,6 @@
<ProjectReference Include="..\..\Framework\Core\Core.csproj" />
<ProjectReference Include="..\..\Framework\KubernetesWorkflow\KubernetesWorkflow.csproj" />
<ProjectReference Include="..\..\Framework\OverwatchTranscript\OverwatchTranscript.csproj" />
<ProjectReference Include="..\CodexClient\CodexClient.csproj" />
<ProjectReference Include="..\CodexContractsPlugin\CodexContractsPlugin.csproj" />
<ProjectReference Include="..\GethPlugin\GethPlugin.csproj" />
<ProjectReference Include="..\MetricsPlugin\MetricsPlugin.csproj" />
@@ -1,160 +0,0 @@
using System.Net.Sockets;
using System.Net;
using Nethereum.Util;
namespace CodexPlugin
{
public class ProcessRecipe
{
public ProcessRecipe(string cmd, string[] args)
{
Cmd = cmd;
Args = args;
}
public string Cmd { get; }
public string[] Args { get; }
}
public class CodexProcessConfig
{
public CodexProcessConfig(string name, FreePortFinder freePortFinder, string dataDir)
{
ApiPort = freePortFinder.GetNextFreePort();
DiscPort = freePortFinder.GetNextFreePort();
ListenPort = freePortFinder.GetNextFreePort();
Name = name;
DataDir = dataDir;
var host = Dns.GetHostEntry(Dns.GetHostName());
var addrs = host.AddressList.Where(a => a.AddressFamily == AddressFamily.InterNetwork).ToList();
LocalIpAddrs = addrs.First();
}
public int ApiPort { get; }
public int DiscPort { get; }
public int ListenPort { get; }
public string Name { get; }
public string DataDir { get; }
public IPAddress LocalIpAddrs { get; }
}
public class CodexProcessRecipe
{
private readonly CodexProcessConfig pc;
private readonly CodexExePath codexExePath;
public CodexProcessRecipe(CodexProcessConfig pc, CodexExePath codexExePath)
{
this.pc = pc;
this.codexExePath = codexExePath;
}
public ProcessRecipe Initialize(CodexStartupConfig config)
{
args.Clear();
AddArg("--api-port", pc.ApiPort);
AddArg("--api-bindaddr", "0.0.0.0");
AddArg("--data-dir", pc.DataDir);
AddArg("--disc-port", pc.DiscPort);
AddArg("--log-level", config.LogLevelWithTopics());
// This makes the node announce itself to its local IP address.
AddArg("--nat", $"extip:{pc.LocalIpAddrs.ToStringInvariant()}");
AddArg("--listen-addrs", $"/ip4/0.0.0.0/tcp/{pc.ListenPort}");
if (!string.IsNullOrEmpty(config.BootstrapSpr))
{
AddArg("--bootstrap-node", config.BootstrapSpr);
}
if (config.StorageQuota != null)
{
AddArg("--storage-quota", config.StorageQuota.SizeInBytes.ToString()!);
}
if (config.BlockTTL != null)
{
AddArg("--block-ttl", config.BlockTTL.ToString()!);
}
if (config.BlockMaintenanceInterval != null)
{
AddArg("--block-mi", Convert.ToInt32(config.BlockMaintenanceInterval.Value.TotalSeconds).ToString());
}
if (config.BlockMaintenanceNumber != null)
{
AddArg("--block-mn", config.BlockMaintenanceNumber.ToString()!);
}
if (config.MetricsEnabled)
{
throw new Exception("Not supported");
//var metricsPort = CreateApiPort(config, MetricsPortTag);
//AddEnvVar("CODEX_METRICS", "true");
//AddEnvVar("CODEX_METRICS_ADDRESS", "0.0.0.0");
//AddEnvVar("CODEX_METRICS_PORT", metricsPort);
//AddPodAnnotation("prometheus.io/scrape", "true");
//AddPodAnnotation("prometheus.io/port", metricsPort.Number.ToString());
}
if (config.SimulateProofFailures != null)
{
throw new Exception("Not supported");
//AddEnvVar("CODEX_SIMULATE_PROOF_FAILURES", config.SimulateProofFailures.ToString()!);
}
if (config.MarketplaceConfig != null)
{
throw new Exception("Not supported");
//var mconfig = config.MarketplaceConfig;
//var gethStart = mconfig.GethNode.StartResult;
//var wsAddress = gethStart.Container.GetInternalAddress(GethContainerRecipe.WsPortTag);
//var marketplaceAddress = mconfig.CodexContracts.Deployment.MarketplaceAddress;
//AddEnvVar("CODEX_ETH_PROVIDER", $"{wsAddress.Host.Replace("http://", "ws://")}:{wsAddress.Port}");
//AddEnvVar("CODEX_MARKETPLACE_ADDRESS", marketplaceAddress);
//var marketplaceSetup = config.MarketplaceConfig.MarketplaceSetup;
//// Custom scripting in the Codex test image will write this variable to a private-key file,
//// and pass the correct filename to Codex.
//var account = marketplaceSetup.EthAccountSetup.GetNew();
//AddEnvVar("ETH_PRIVATE_KEY", account.PrivateKey);
//Additional(account);
//SetCommandOverride(marketplaceSetup);
//if (marketplaceSetup.IsValidator)
//{
// AddEnvVar("CODEX_VALIDATOR", "true");
//}
}
//if (!string.IsNullOrEmpty(config.NameOverride))
//{
// AddEnvVar("CODEX_NODENAME", config.NameOverride);
//}
return Create();
}
private ProcessRecipe Create()
{
return new ProcessRecipe(
cmd: codexExePath.Get(),
args: args.ToArray());
}
private readonly List<string> args = new List<string>();
private void AddArg(string arg, string val)
{
args.Add($"{arg}={val}");
}
private void AddArg(string arg, int val)
{
args.Add($"{arg}={val}");
}
}
}
+2 -4
View File
@@ -1,5 +1,4 @@
using CodexClient;
using CodexContractsPlugin;
using CodexContractsPlugin;
using GethPlugin;
using KubernetesWorkflow;
using Utils;
@@ -54,7 +53,6 @@ namespace CodexPlugin
public CodexLogLevel ContractClock { get; set; } = CodexLogLevel.Warn;
public CodexLogLevel? BlockExchange { get; }
public CodexLogLevel JsonSerialize { get; set; } = CodexLogLevel.Warn;
public CodexLogLevel MarketplaceInfra { get; set; } = CodexLogLevel.Warn;
}
public class CodexSetup : CodexStartupConfig, ICodexSetup
@@ -225,7 +223,7 @@ namespace CodexPlugin
{
if (pinned) return accounts.Last();
var a = EthAccountGenerator.GenerateNew();
var a = EthAccount.GenerateNew();
accounts.Add(a);
return a;
}
@@ -1,26 +1,29 @@
using CodexClient;
using CodexPlugin.Hooks;
using Core;
using GethPlugin;
using KubernetesWorkflow;
using KubernetesWorkflow.Types;
using Utils;
using Logging;
namespace CodexPlugin
{
public class ContainerCodexStarter : ICodexStarter
public class CodexStarter
{
private readonly IPluginTools pluginTools;
private readonly ProcessControlMap processControlMap;
private readonly CodexContainerRecipe recipe = new CodexContainerRecipe();
private readonly ApiChecker apiChecker;
private DebugInfoVersion? versionResponse;
public ContainerCodexStarter(IPluginTools pluginTools, ProcessControlMap processControlMap)
public CodexStarter(IPluginTools pluginTools)
{
this.pluginTools = pluginTools;
this.processControlMap = processControlMap;
apiChecker = new ApiChecker(pluginTools);
}
public ICodexInstance[] BringOnline(CodexSetup codexSetup)
public CodexHooksFactory HooksFactory { get; } = new CodexHooksFactory();
public RunningPod[] BringOnline(CodexSetup codexSetup)
{
LogSeparator();
Log($"Starting {codexSetup.Describe()}...");
@@ -40,11 +43,51 @@ namespace CodexPlugin
}
LogSeparator();
return containers.Select(CreateInstance).ToArray();
return containers;
}
public void Decommission()
public ICodexNodeGroup WrapCodexContainers(CoreInterface coreInterface, RunningPod[] containers)
{
var codexNodeFactory = new CodexNodeFactory(pluginTools, HooksFactory);
var group = CreateCodexGroup(coreInterface, containers, codexNodeFactory);
Log($"Codex version: {group.Version}");
versionResponse = group.Version;
return group;
}
public void BringOffline(CodexNodeGroup group, bool waitTillStopped)
{
Log($"Stopping {group.Describe()}...");
StopCrashWatcher(group);
var workflow = pluginTools.CreateWorkflow();
foreach (var c in group.Containers)
{
workflow.Stop(c, waitTillStopped);
}
Log("Stopped.");
}
public void Stop(RunningPod pod, bool waitTillStopped)
{
Log($"Stopping node...");
var workflow = pluginTools.CreateWorkflow();
workflow.Stop(pod, waitTillStopped);
Log("Stopped.");
}
public string GetCodexId()
{
if (versionResponse != null) return versionResponse.Version;
return recipe.Image;
}
public string GetCodexRevision()
{
if (versionResponse != null) return versionResponse.Revision;
return "unknown";
}
private StartupConfig CreateStartupConfig(CodexSetup codexSetup)
@@ -75,15 +118,27 @@ namespace CodexPlugin
return workflow.GetPodInfo(rc);
}
private ICodexInstance CreateInstance(RunningPod pod)
private CodexNodeGroup CreateCodexGroup(CoreInterface coreInterface, RunningPod[] runningContainers, CodexNodeFactory codexNodeFactory)
{
var instance = CodexInstanceContainerExtension.CreateFromPod(pod);
var processControl = new CodexContainerProcessControl(pluginTools, pod, onStop: () =>
var group = new CodexNodeGroup(this, pluginTools, runningContainers, codexNodeFactory);
try
{
processControlMap.Remove(instance);
});
processControlMap.Add(instance, processControl);
return instance;
Stopwatch.Measure(pluginTools.GetLog(), "EnsureOnline", group.EnsureOnline);
}
catch
{
CodexNodesNotOnline(coreInterface, runningContainers);
throw;
}
return group;
}
private void CodexNodesNotOnline(CoreInterface coreInterface, RunningPod[] runningContainers)
{
Log("Codex nodes failed to start");
foreach (var container in runningContainers.First().Containers) coreInterface.DownloadLog(container);
}
private void LogSeparator()
@@ -102,5 +157,13 @@ namespace CodexPlugin
{
pluginTools.GetLog().Log(message);
}
private void StopCrashWatcher(CodexNodeGroup group)
{
foreach (var node in group)
{
node.CrashWatcher.Stop();
}
}
}
}
@@ -1,5 +1,4 @@
using CodexClient;
using KubernetesWorkflow;
using KubernetesWorkflow;
using Utils;
namespace CodexPlugin
@@ -30,7 +29,6 @@ namespace CodexPlugin
{
"discv5",
"providers",
"routingtable",
"manager",
"cache",
};
@@ -81,20 +79,12 @@ namespace CodexPlugin
"json",
"serialization"
};
var marketplaceInfraTopics = new[]
{
"JSONRPC-WS-CLIENT",
"JSONRPC-HTTP-CLIENT",
"codex",
"repostore"
};
level = $"{level};" +
$"{CustomTopics.DiscV5.ToString()!.ToLowerInvariant()}:{string.Join(",", discV5Topics)};" +
$"{CustomTopics.Libp2p.ToString()!.ToLowerInvariant()}:{string.Join(",", libp2pTopics)};" +
$"{CustomTopics.ContractClock.ToString().ToLowerInvariant()}:{string.Join(",", contractClockTopics)};" +
$"{CustomTopics.JsonSerialize.ToString().ToLowerInvariant()}:{string.Join(",", jsonSerializeTopics)};" +
$"{CustomTopics.MarketplaceInfra.ToString().ToLowerInvariant()}:{string.Join(",", marketplaceInfraTopics)}";
$"{CustomTopics.JsonSerialize.ToString().ToLowerInvariant()}:{string.Join(",", jsonSerializeTopics)}";
if (CustomTopics.BlockExchange != null)
{
@@ -1,7 +1,7 @@
using Newtonsoft.Json;
using Utils;
namespace CodexClient
namespace CodexPlugin
{
public class DebugInfo
{
@@ -109,16 +109,6 @@ namespace CodexClient
{
return HashCode.Combine(Id);
}
public static bool operator ==(ContentId a, ContentId b)
{
return a.Id == b.Id;
}
public static bool operator !=(ContentId a, ContentId b)
{
return a.Id != b.Id;
}
}
public class CodexSpace
@@ -1,4 +1,4 @@
namespace CodexClient
namespace CodexPlugin
{
public static class CodexUtils
{
@@ -1,80 +0,0 @@
using CodexClient;
using CodexClient.Hooks;
using Core;
using Logging;
namespace CodexPlugin
{
public class CodexWrapper
{
private readonly IPluginTools pluginTools;
private readonly ProcessControlMap processControlMap;
private readonly CodexHooksFactory hooksFactory;
private DebugInfoVersion? versionResponse;
public CodexWrapper(IPluginTools pluginTools, ProcessControlMap processControlMap, CodexHooksFactory hooksFactory)
{
this.pluginTools = pluginTools;
this.processControlMap = processControlMap;
this.hooksFactory = hooksFactory;
}
public string GetCodexId()
{
if (versionResponse != null) return versionResponse.Version;
return "unknown";
}
public string GetCodexRevision()
{
if (versionResponse != null) return versionResponse.Revision;
return "unknown";
}
public ICodexNodeGroup WrapCodexInstances(ICodexInstance[] instances)
{
var codexNodeFactory = new CodexNodeFactory(
log: pluginTools.GetLog(),
fileManager: pluginTools.GetFileManager(),
hooksFactory: hooksFactory,
httpFactory: pluginTools,
processControlFactory: processControlMap);
var group = CreateCodexGroup(instances, codexNodeFactory);
pluginTools.GetLog().Log($"Codex version: {group.Version}");
versionResponse = group.Version;
return group;
}
private CodexNodeGroup CreateCodexGroup(ICodexInstance[] instances, CodexNodeFactory codexNodeFactory)
{
var nodes = instances.Select(codexNodeFactory.CreateCodexNode).ToArray();
var group = new CodexNodeGroup(pluginTools, nodes);
try
{
Stopwatch.Measure(pluginTools.GetLog(), "EnsureOnline", group.EnsureOnline);
}
catch
{
CodexNodesNotOnline(instances);
throw;
}
return group;
}
private void CodexNodesNotOnline(ICodexInstance[] instances)
{
pluginTools.GetLog().Log("Codex nodes failed to start");
var log = pluginTools.GetLog();
foreach (var i in instances)
{
var pc = processControlMap.Get(i);
pc.DownloadLog(log.CreateSubfile(i.Name + "_failed_to_start"));
}
}
}
}
@@ -1,19 +1,19 @@
using CodexClient;
using CodexClient.Hooks;
using CodexPlugin.Hooks;
using Core;
using KubernetesWorkflow.Types;
namespace CodexPlugin
{
public static class CoreInterfaceExtensions
{
public static ICodexInstance[] DeployCodexNodes(this CoreInterface ci, int number, Action<ICodexSetup> setup)
public static RunningPod[] DeployCodexNodes(this CoreInterface ci, int number, Action<ICodexSetup> setup)
{
return Plugin(ci).DeployCodexNodes(number, setup);
}
public static ICodexNodeGroup WrapCodexContainers(this CoreInterface ci, ICodexInstance[] instances)
public static ICodexNodeGroup WrapCodexContainers(this CoreInterface ci, RunningPod[] containers)
{
return Plugin(ci).WrapCodexContainers(instances);
return Plugin(ci).WrapCodexContainers(ci, containers);
}
public static ICodexNode StartCodexNode(this CoreInterface ci)
@@ -39,9 +39,9 @@ namespace CodexPlugin
return ci.StartCodexNodes(number, s => { });
}
public static void AddCodexHooksProvider(this CoreInterface ci, ICodexHooksProvider hooksProvider)
public static void SetCodexHooksProvider(this CoreInterface ci, ICodexHooksProvider hooksProvider)
{
Plugin(ci).AddCodexHooksProvider(hooksProvider);
Plugin(ci).SetCodexHooksProvider(hooksProvider);
}
private static CodexPlugin Plugin(CoreInterface ci)
@@ -1,43 +0,0 @@
using System.Net.NetworkInformation;
namespace CodexPlugin
{
public class FreePortFinder
{
private readonly object _lock = new object();
private int nextPort = 8080;
public int GetNextFreePort()
{
lock (_lock)
{
return Next();
}
}
private int Next()
{
while (true)
{
var p = nextPort;
nextPort++;
if (!IsInUse(p))
{
return p;
}
if (nextPort > 30000) throw new Exception("Running out of ports.");
}
}
private bool IsInUse(int port)
{
var ipProps = IPGlobalProperties.GetIPGlobalProperties();
if (ipProps.GetActiveTcpConnections().Any(t => t.LocalEndPoint.Port == port)) return true;
if (ipProps.GetActiveTcpListeners().Any(t => t.Port == port)) return true;
if (ipProps.GetActiveUdpListeners().Any(u => u.Port == port)) return true;
return false;
}
}
}
@@ -1,6 +1,7 @@
using Utils;
using GethPlugin;
using Utils;
namespace CodexClient.Hooks
namespace CodexPlugin.Hooks
{
public interface ICodexHooksProvider
{
@@ -9,14 +10,11 @@ namespace CodexClient.Hooks
public class CodexHooksFactory
{
public List<ICodexHooksProvider> Providers { get; } = new List<ICodexHooksProvider>();
public ICodexHooksProvider Provider { get; set; } = new DoNothingHooksProvider();
public ICodexNodeHooks CreateHooks(string nodeName)
{
if (Providers.Count == 0) return new DoNothingCodexHooks();
var hooks = Providers.Select(p => p.CreateHooks(nodeName)).ToArray();
return new MuxingCodexNodeHooks(hooks);
return Provider.CreateHooks(nodeName);
}
}
@@ -46,7 +44,7 @@ namespace CodexClient.Hooks
{
}
public void OnNodeStarted(ICodexNode node, string peerId, string nodeId)
public void OnNodeStarted(string peerId, string nodeId)
{
}
@@ -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,10 +0,0 @@
using CodexClient;
namespace CodexPlugin
{
public interface ICodexStarter
{
ICodexInstance[] BringOnline(CodexSetup codexSetup);
void Decommission();
}
}
@@ -1,9 +1,10 @@
using CodexOpenApi;
using CodexContractsPlugin;
using CodexOpenApi;
using Newtonsoft.Json.Linq;
using System.Numerics;
using Utils;
namespace CodexClient
namespace CodexPlugin
{
public class Mapper
{
@@ -14,7 +15,7 @@ namespace CodexClient
Id = debugInfo.Id,
Spr = debugInfo.Spr,
Addrs = debugInfo.Addrs.ToArray(),
AnnounceAddresses = debugInfo.AnnounceAddresses.ToArray(),
AnnounceAddresses = JArray(debugInfo.AdditionalProperties, "announceAddresses").Select(x => x.ToString()).ToArray(),
Version = Map(debugInfo.Codex),
Table = Map(debugInfo.Table)
};
@@ -42,8 +43,8 @@ namespace CodexClient
return new CodexOpenApi.SalesAvailabilityCREATE
{
Duration = ToDecInt(availability.MaxDuration.TotalSeconds),
MinPricePerBytePerSecond = ToDecInt(availability.MinPricePerBytePerSecond),
TotalCollateral = ToDecInt(availability.TotalCollateral),
MinPrice = ToDecInt(availability.MinPriceForTotalSpace),
MaxCollateral = ToDecInt(availability.MaxCollateral),
TotalSize = ToDecInt(availability.TotalSpace.SizeInBytes)
};
}
@@ -54,27 +55,27 @@ namespace CodexClient
{
Duration = ToDecInt(purchase.Duration.TotalSeconds),
ProofProbability = ToDecInt(purchase.ProofProbability),
PricePerBytePerSecond = ToDecInt(purchase.PricePerBytePerSecond),
CollateralPerByte = ToDecInt(purchase.CollateralPerByte),
Reward = ToDecInt(purchase.PricePerSlotPerSecond),
Collateral = ToDecInt(purchase.RequiredCollateral),
Expiry = ToDecInt(purchase.Expiry.TotalSeconds),
Nodes = Convert.ToInt32(purchase.MinRequiredNumberOfNodes),
Tolerance = Convert.ToInt32(purchase.NodeFailureTolerance)
};
}
public StorageAvailability[] Map(ICollection<CodexOpenApi.SalesAvailabilityREAD> availabilities)
public StorageAvailability[] Map(ICollection<SalesAvailabilityREAD> availabilities)
{
return availabilities.Select(a => Map(a)).ToArray();
}
public StorageAvailability Map(CodexOpenApi.SalesAvailabilityREAD availability)
public StorageAvailability Map(SalesAvailabilityREAD availability)
{
return new StorageAvailability
(
ToByteSize(availability.TotalSize),
ToTimespan(availability.Duration),
new TestToken(ToBigIng(availability.MinPricePerBytePerSecond)),
new TestToken(ToBigIng(availability.TotalCollateral))
new TestToken(ToBigIng(availability.MinPrice)),
new TestToken(ToBigIng(availability.MaxCollateral))
)
{
Id = availability.Id,
@@ -82,80 +83,49 @@ namespace CodexClient
};
}
public StoragePurchase Map(CodexOpenApi.Purchase purchase)
{
return new StoragePurchase
{
Request = Map(purchase.Request),
State = purchase.State.ToString(), //Map(purchase.State),
Error = purchase.Error
};
}
// TODO: Fix openapi spec for this call.
//public StoragePurchase Map(CodexOpenApi.Purchase purchase)
//{
// return new StoragePurchase(Map(purchase.Request))
// {
// State = purchase.State,
// Error = purchase.Error
// };
//}
public StoragePurchaseState Map(PurchaseState purchaseState)
{
// TODO: to be re-enabled when marketplace api lines up with openapi.yaml.
//public StorageRequest Map(CodexOpenApi.StorageRequest request)
//{
// return new StorageRequest(Map(request.Ask), Map(request.Content))
// {
// Id = request.Id,
// Client = request.Client,
// Expiry = TimeSpan.FromSeconds(Convert.ToInt64(request.Expiry)),
// Nonce = request.Nonce
// };
//}
// Explicit mapping: If the API changes, we will get compile errors here.
// That's what we want.
switch (purchaseState)
{
case PurchaseState.Cancelled:
return StoragePurchaseState.Cancelled;
case PurchaseState.Error:
return StoragePurchaseState.Error;
case PurchaseState.Failed:
return StoragePurchaseState.Failed;
case PurchaseState.Finished:
return StoragePurchaseState.Finished;
case PurchaseState.Pending:
return StoragePurchaseState.Pending;
case PurchaseState.Started:
return StoragePurchaseState.Started;
case PurchaseState.Submitted:
return StoragePurchaseState.Submitted;
case PurchaseState.Unknown:
return StoragePurchaseState.Unknown;
}
//public StorageAsk Map(CodexOpenApi.StorageAsk ask)
//{
// return new StorageAsk
// {
// Duration = TimeSpan.FromSeconds(Convert.ToInt64(ask.Duration)),
// MaxSlotLoss = ask.MaxSlotLoss,
// ProofProbability = ask.ProofProbability,
// Reward = Convert.ToDecimal(ask.Reward).TstWei(),
// Slots = ask.Slots,
// SlotSize = new ByteSize(Convert.ToInt64(ask.SlotSize))
// };
//}
throw new Exception("API incompatibility detected. Unknown purchaseState: " + purchaseState.ToString());
}
//public StorageContent Map(CodexOpenApi.Content content)
//{
// return new StorageContent
// {
// Cid = content.Cid
// };
//}
public StorageRequest Map(CodexOpenApi.StorageRequest request)
{
return new StorageRequest
{
Ask = Map(request.Ask),
Content = Map(request.Content),
Id = request.Id,
Client = request.Client,
Expiry = request.Expiry,
Nonce = request.Nonce
};
}
public StorageAsk Map(CodexOpenApi.StorageAsk ask)
{
return new StorageAsk
{
Duration = ask.Duration,
MaxSlotLoss = ask.MaxSlotLoss,
ProofProbability = ask.ProofProbability,
PricePerBytePerSecond = ask.PricePerBytePerSecond,
Slots = ask.Slots,
SlotSize = ask.SlotSize
};
}
public StorageContent Map(CodexOpenApi.Content content)
{
return new StorageContent
{
Cid = content.Cid
};
}
public CodexSpace Map(CodexOpenApi.Space space)
public CodexSpace Map(Space space)
{
return new CodexSpace
{
@@ -166,7 +136,7 @@ namespace CodexClient
};
}
private DebugInfoVersion Map(CodexOpenApi.CodexVersion obj)
private DebugInfoVersion Map(CodexVersion obj)
{
return new DebugInfoVersion
{
@@ -175,7 +145,7 @@ namespace CodexClient
};
}
private DebugInfoTable Map(CodexOpenApi.PeersTable obj)
private DebugInfoTable Map(PeersTable obj)
{
return new DebugInfoTable
{
@@ -184,7 +154,7 @@ namespace CodexClient
};
}
private DebugInfoTableNode Map(CodexOpenApi.Node? token)
private DebugInfoTableNode Map(Node? token)
{
if (token == null) return new DebugInfoTableNode();
return new DebugInfoTableNode
@@ -197,7 +167,7 @@ namespace CodexClient
};
}
private DebugInfoTableNode[] Map(ICollection<CodexOpenApi.Node> nodes)
private DebugInfoTableNode[] Map(ICollection<Node> nodes)
{
if (nodes == null || nodes.Count == 0)
{
@@ -1,8 +1,8 @@
using CodexClient.Hooks;
using CodexPlugin.Hooks;
using Logging;
using Utils;
namespace CodexClient
namespace CodexPlugin
{
public interface IMarketplaceAccess
{
@@ -27,12 +27,8 @@ namespace CodexClient
public IStoragePurchaseContract RequestStorage(StoragePurchaseRequest purchase)
{
purchase.Log(log);
var swResult = Stopwatch.Measure(log, nameof(RequestStorage), () =>
{
return codexAccess.RequestStorage(purchase);
});
var response = swResult.Value;
var response = codexAccess.RequestStorage(purchase);
if (string.IsNullOrEmpty(response) ||
response == "Unable to encode manifest" ||
@@ -46,7 +42,12 @@ namespace CodexClient
Log($"Storage requested successfully. PurchaseId: '{response}'.");
return new StoragePurchaseContract(log, codexAccess, response, purchase, hooks);
var contract = new StoragePurchaseContract(log, codexAccess, response, purchase, hooks);
contract.WaitForStorageContractSubmitted();
hooks.OnStorageContractSubmitted(contract);
return contract;
}
public string MakeStorageAvailable(StorageAvailability availability)
@@ -71,7 +72,7 @@ namespace CodexClient
private void Log(string msg)
{
log.Log($"{codexAccess.GetName()} {msg}");
log.Log($"{codexAccess.Container.Containers.Single().Name} {msg}");
}
}
@@ -1,7 +1,8 @@
using Logging;
using CodexContractsPlugin;
using Logging;
using Utils;
namespace CodexClient
namespace CodexPlugin
{
public class StoragePurchaseRequest
{
@@ -11,8 +12,8 @@ namespace CodexClient
}
public ContentId ContentId { get; set; }
public TestToken PricePerBytePerSecond { get; set; } = 1.TstWei();
public TestToken CollateralPerByte { get; set; } = 1.TstWei();
public TestToken PricePerSlotPerSecond { get; set; } = 1.TstWei();
public TestToken RequiredCollateral { get; set; } = 1.TstWei();
public uint MinRequiredNumberOfNodes { get; set; }
public uint NodeFailureTolerance { get; set; }
public int ProofProbability { get; set; }
@@ -22,8 +23,8 @@ namespace CodexClient
public void Log(ILog log)
{
log.Log($"Requesting storage for: {ContentId.Id}... (" +
$"pricePerBytePerSecond: {PricePerBytePerSecond}, " +
$"collateralPerByte: {CollateralPerByte}, " +
$"pricePerSlotPerSecond: {PricePerSlotPerSecond}, " +
$"requiredCollateral: {RequiredCollateral}, " +
$"minRequiredNumberOfNodes: {MinRequiredNumberOfNodes}, " +
$"nodeFailureTolerance: {NodeFailureTolerance}, " +
$"proofProbability: {ProofProbability}, " +
@@ -37,24 +38,6 @@ namespace CodexClient
public string State { get; set; } = string.Empty;
public string Error { get; set; } = string.Empty;
public StorageRequest Request { get; set; } = null!;
public bool IsCancelled => State.ToLowerInvariant().Contains("cancel");
public bool IsError => State.ToLowerInvariant().Contains("error");
public bool IsFinished => State.ToLowerInvariant().Contains("finished");
public bool IsStarted => State.ToLowerInvariant().Contains("started");
public bool IsSubmitted => State.ToLowerInvariant().Contains("submitted");
}
public enum StoragePurchaseState
{
Cancelled = 0,
Error = 1,
Failed = 2,
Finished = 3,
Pending = 4,
Started = 5,
Submitted = 6,
Unknown = 7,
}
public class StorageRequest
@@ -73,7 +56,7 @@ namespace CodexClient
public string SlotSize { get; set; } = string.Empty;
public string Duration { get; set; } = string.Empty;
public string ProofProbability { get; set; } = string.Empty;
public string PricePerBytePerSecond { get; set; } = string.Empty;
public string Reward { get; set; } = string.Empty;
public int MaxSlotLoss { get; set; }
}
@@ -86,19 +69,19 @@ namespace CodexClient
public class StorageAvailability
{
public StorageAvailability(ByteSize totalSpace, TimeSpan maxDuration, TestToken minPricePerBytePerSecond, TestToken totalCollateral)
public StorageAvailability(ByteSize totalSpace, TimeSpan maxDuration, TestToken minPriceForTotalSpace, TestToken maxCollateral)
{
TotalSpace = totalSpace;
MaxDuration = maxDuration;
MinPricePerBytePerSecond = minPricePerBytePerSecond;
TotalCollateral = totalCollateral;
MinPriceForTotalSpace = minPriceForTotalSpace;
MaxCollateral = maxCollateral;
}
public string Id { get; set; } = string.Empty;
public ByteSize TotalSpace { get; }
public TimeSpan MaxDuration { get; }
public TestToken MinPricePerBytePerSecond { get; }
public TestToken TotalCollateral { get; }
public TestToken MinPriceForTotalSpace { get; }
public TestToken MaxCollateral { get; }
public ByteSize FreeSpace { get; set; } = ByteSize.Zero;
public void Log(ILog log)
@@ -106,8 +89,8 @@ namespace CodexClient
log.Log($"Storage Availability: (" +
$"totalSize: {TotalSpace}, " +
$"maxDuration: {Time.FormatDuration(MaxDuration)}, " +
$"minPricePerBytePerSecond: {MinPricePerBytePerSecond}, " +
$"totalCollateral: {TotalCollateral})");
$"minPriceForTotalSpace: {MinPriceForTotalSpace}, " +
$"maxCollateral: {MaxCollateral})");
}
}
}
@@ -1,6 +1,5 @@
using CodexClient;
using CodexPlugin.OverwatchSupport.LineConverters;
using Logging;
using CodexPlugin.OverwatchSupport.LineConverters;
using KubernetesWorkflow;
using OverwatchTranscript;
using Utils;
@@ -1,5 +1,5 @@
using CodexClient;
using CodexClient.Hooks;
using CodexPlugin.Hooks;
using GethPlugin;
using OverwatchTranscript;
using Utils;
@@ -32,7 +32,7 @@ namespace CodexPlugin.OverwatchSupport
});
}
public void OnNodeStarted(ICodexNode node, string peerId, string nodeId)
public void OnNodeStarted(string peerId, string nodeId)
{
if (string.IsNullOrEmpty(peerId) || string.IsNullOrEmpty(nodeId))
{
@@ -1,4 +1,5 @@
using CodexClient.Hooks;
using CodexPlugin.Hooks;
using KubernetesWorkflow;
using Logging;
using OverwatchTranscript;
using Utils;

Some files were not shown because too many files have changed in this diff Show More