2023-09-12 13:32:06 +02:00
using Core ;
2023-09-11 10:43:27 +02:00
using FileUtils ;
2023-09-19 11:51:59 +02:00
using GethPlugin ;
2023-09-13 12:24:46 +02:00
using KubernetesWorkflow ;
2023-11-12 10:07:23 +01:00
using KubernetesWorkflow.Types ;
2023-09-12 11:43:46 +02:00
using Logging ;
2023-09-13 11:59:21 +02:00
using MetricsPlugin ;
2023-08-02 15:11:27 +02:00
using Utils ;
2023-04-13 09:33:10 +02:00
2023-09-11 11:59:33 +02:00
namespace CodexPlugin
2023-04-13 09:33:10 +02:00
{
2023-09-19 11:51:59 +02:00
public interface ICodexNode : IHasContainer , IHasMetricsScrapeTarget , IHasEthAddress
2023-04-13 09:33:10 +02:00
{
2023-04-19 09:19:06 +02:00
string GetName ();
2023-04-13 09:33:10 +02:00
CodexDebugResponse GetDebugInfo ();
2023-05-11 12:44:53 +02:00
CodexDebugPeerResponse GetDebugPeer ( string peerId );
2023-10-23 09:36:31 +02:00
// These debug methods are not available in master-line Codex. Use only for custom builds.
//CodexDebugBlockExchangeResponse GetDebugBlockExchange();
//CodexDebugRepoStoreResponse[] GetDebugRepoStore();
2023-09-12 13:32:06 +02:00
ContentId UploadFile ( TrackedFile file );
TrackedFile ? DownloadContent ( ContentId contentId , string fileLabel = "" );
2023-11-10 08:20:08 +01:00
CodexLocalData [] LocalFiles ();
2023-09-19 11:51:59 +02:00
void ConnectToPeer ( ICodexNode node );
2023-07-31 11:51:29 +02:00
CodexDebugVersionResponse Version { get ; }
2023-09-20 08:45:55 +02:00
IMarketplaceAccess Marketplace { get ; }
2023-09-20 12:55:09 +02:00
CrashWatcher CrashWatcher { get ; }
2023-11-06 14:33:47 +01:00
PodInfo GetPodInfo ();
2023-09-19 11:51:59 +02:00
void Stop ();
2023-04-13 09:33:10 +02:00
}
2023-09-19 11:51:59 +02:00
public class CodexNode : ICodexNode
2023-04-13 09:33:10 +02:00
{
private const string SuccessfullyConnectedMessage = "Successfully connected to peer" ;
private const string UploadFailedMessage = "Unable to store block" ;
2023-09-12 11:43:46 +02:00
private readonly IPluginTools tools ;
2023-09-20 10:13:29 +02:00
private readonly EthAddress ? ethAddress ;
2023-04-13 09:33:10 +02:00
2023-09-20 10:13:29 +02:00
public CodexNode ( IPluginTools tools , CodexAccess codexAccess , CodexNodeGroup group , IMarketplaceAccess marketplaceAccess , EthAddress ? ethAddress )
2023-04-13 09:33:10 +02:00
{
2023-09-12 11:43:46 +02:00
this . tools = tools ;
2023-09-19 11:51:59 +02:00
this . ethAddress = ethAddress ;
2023-04-13 09:33:10 +02:00
CodexAccess = codexAccess ;
Group = group ;
2023-09-20 08:45:55 +02:00
Marketplace = marketplaceAccess ;
2023-07-31 11:51:29 +02:00
Version = new CodexDebugVersionResponse ();
2023-04-13 09:33:10 +02:00
}
2023-09-13 12:24:46 +02:00
public RunningContainer Container { get { return CodexAccess . Container ; } }
2023-04-13 09:33:10 +02:00
public CodexAccess CodexAccess { get ; }
2023-09-20 12:55:09 +02:00
public CrashWatcher CrashWatcher { get => CodexAccess . CrashWatcher ; }
2023-04-13 09:33:10 +02:00
public CodexNodeGroup Group { get ; }
2023-09-20 08:45:55 +02:00
public IMarketplaceAccess Marketplace { get ; }
2023-07-31 11:51:29 +02:00
public CodexDebugVersionResponse Version { get ; private set ; }
2023-09-13 11:59:21 +02:00
public IMetricsScrapeTarget MetricsScrapeTarget
{
get
{
2023-11-06 14:33:47 +01:00
return new MetricsScrapeTarget ( CodexAccess . Container , CodexContainerRecipe . MetricsPortTag );
2023-09-13 11:59:21 +02:00
}
}
2023-09-20 10:13:29 +02:00
public EthAddress EthAddress
2023-09-19 11:51:59 +02:00
{
get
{
if ( ethAddress == null ) throw new Exception ( "Marketplace is not enabled for this Codex node. Please start it with the option '.EnableMarketplace(...)' to enable it." );
return ethAddress ;
}
}
2023-04-13 09:33:10 +02:00
public string GetName ()
{
2023-04-30 10:56:19 +02:00
return CodexAccess . Container . Name ;
2023-04-13 09:33:10 +02:00
}
public CodexDebugResponse GetDebugInfo ()
{
2023-06-29 16:07:49 +02:00
var debugInfo = CodexAccess . GetDebugInfo ();
2023-05-18 10:42:04 +02:00
var known = string . Join ( "," , debugInfo . table . nodes . Select ( n => n . peerId ));
Log ( $"Got DebugInfo with id: '{debugInfo.id}'. This node knows: {known}" );
2023-04-19 14:57:00 +02:00
return debugInfo ;
2023-04-13 09:33:10 +02:00
}
2023-05-11 12:44:53 +02:00
public CodexDebugPeerResponse GetDebugPeer ( string peerId )
{
2023-06-29 16:07:49 +02:00
return CodexAccess . GetDebugPeer ( peerId );
2023-05-11 12:44:53 +02:00
}
2023-10-09 12:33:35 +02:00
public CodexDebugBlockExchangeResponse GetDebugBlockExchange ()
{
return CodexAccess . GetDebugBlockExchange ();
}
2023-10-10 17:54:19 +02:00
public CodexDebugRepoStoreResponse [] GetDebugRepoStore ()
{
return CodexAccess . GetDebugRepoStore ();
}
2023-09-12 13:32:06 +02:00
public ContentId UploadFile ( TrackedFile file )
2023-04-13 09:33:10 +02:00
{
2023-09-12 11:43:46 +02:00
using var fileStream = File . OpenRead ( file . Filename );
2023-06-07 08:30:10 +02:00
2023-09-12 11:43:46 +02:00
var logMessage = $"Uploading file {file.Describe()}..." ;
Log ( logMessage );
var response = Stopwatch . Measure ( tools . GetLog (), logMessage , () =>
{
return CodexAccess . UploadFile ( fileStream );
});
2023-06-07 08:30:10 +02:00
2023-09-20 10:51:47 +02:00
if ( string . IsNullOrEmpty ( response )) FrameworkAssert . Fail ( "Received empty response." );
if ( response . StartsWith ( UploadFailedMessage )) FrameworkAssert . Fail ( "Node failed to store block." );
2023-07-11 08:19:14 +02:00
2023-09-12 11:43:46 +02:00
Log ( $"Uploaded file. Received contentId: '{response}'." );
return new ContentId ( response );
2023-04-13 09:33:10 +02:00
}
2023-09-12 13:32:06 +02:00
public TrackedFile ? DownloadContent ( ContentId contentId , string fileLabel = "" )
2023-04-13 09:33:10 +02:00
{
2023-09-12 11:43:46 +02:00
var logMessage = $"Downloading for contentId: '{contentId.Id}'..." ;
Log ( logMessage );
2023-09-13 10:03:11 +02:00
var file = tools . GetFileManager (). CreateEmptyFile ( fileLabel );
2023-09-12 11:43:46 +02:00
Stopwatch . Measure ( tools . GetLog (), logMessage , () => DownloadToFile ( contentId . Id , file ));
Log ( $"Downloaded file {file.Describe()} to '{file.Filename}'." );
return file ;
2023-04-13 09:33:10 +02:00
}
2023-11-10 08:20:08 +01:00
public CodexLocalData [] LocalFiles ()
{
return CodexAccess . LocalFiles (). Select ( l => new CodexLocalData ( new ContentId ( l . cid ), l . manifest )). ToArray ();
}
2023-09-19 11:51:59 +02:00
public void ConnectToPeer ( ICodexNode node )
2023-04-13 09:33:10 +02:00
{
2023-09-19 11:51:59 +02:00
var peer = ( CodexNode ) node ;
2023-04-13 09:33:10 +02:00
Log ( $"Connecting to peer {peer.GetName()}..." );
var peerInfo = node . GetDebugInfo ();
2023-06-29 16:07:49 +02:00
var response = CodexAccess . ConnectToPeer ( peerInfo . id , GetPeerMultiAddress ( peer , peerInfo ));
2023-04-13 09:33:10 +02:00
2023-09-20 10:51:47 +02:00
FrameworkAssert . That ( response == SuccessfullyConnectedMessage , "Unable to connect codex nodes." );
2023-04-13 09:33:10 +02:00
Log ( $"Successfully connected to peer {peer.GetName()}." );
}
2023-11-06 14:33:47 +01:00
public PodInfo GetPodInfo ()
{
return CodexAccess . GetPodInfo ();
}
2023-09-19 11:51:59 +02:00
public void Stop ()
2023-04-25 12:52:11 +02:00
{
2023-05-04 14:55:39 +02:00
if ( Group . Count () > 1 ) throw new InvalidOperationException ( "Codex-nodes that are part of a group cannot be " +
"individually shut down. Use 'BringOffline()' on the group object to stop the group. This method is only " +
"available for codex-nodes in groups of 1." );
2023-09-11 16:57:57 +02:00
Group . BringOffline ();
2023-04-25 12:52:11 +02:00
}
2023-07-31 11:51:29 +02:00
public void EnsureOnlineGetVersionResponse ()
{
2023-08-02 15:11:27 +02:00
var debugInfo = Time . Retry ( CodexAccess . GetDebugInfo , "ensure online" );
2023-07-31 11:51:29 +02:00
var nodePeerId = debugInfo . id ;
var nodeName = CodexAccess . Container . Name ;
if (! debugInfo . codex . IsValid ())
{
throw new Exception ( $"Invalid version information received from Codex node {GetName()}: {debugInfo.codex}" );
}
2023-11-12 11:24:58 +01:00
var log = tools . GetLog ();
log . AddStringReplace ( nodePeerId , nodeName );
log . AddStringReplace ( debugInfo . table . localNode . nodeId , nodeName );
2023-07-31 11:51:29 +02:00
Version = debugInfo . codex ;
}
2023-09-19 11:51:59 +02:00
private string GetPeerMultiAddress ( CodexNode peer , CodexDebugResponse peerInfo )
2023-04-13 09:33:10 +02:00
{
var multiAddress = peerInfo . addrs . First ();
// Todo: Is there a case where First address in list is not the way?
// The peer we want to connect is in a different pod.
// We must replace the default IP with the pod IP in the multiAddress.
2023-11-06 14:33:47 +01:00
var workflow = tools . CreateWorkflow ();
var podInfo = workflow . GetPodInfo ( peer . Container );
return multiAddress . Replace ( "0.0.0.0" , podInfo . Ip );
2023-04-13 09:33:10 +02:00
}
2023-09-12 13:32:06 +02:00
private void DownloadToFile ( string contentId , TrackedFile file )
2023-04-13 09:33:10 +02:00
{
using var fileStream = File . OpenWrite ( file . Filename );
2023-04-19 14:57:00 +02:00
try
{
2023-06-29 16:07:49 +02:00
using var downloadStream = CodexAccess . DownloadFile ( contentId );
2023-04-19 14:57:00 +02:00
downloadStream . CopyTo ( fileStream );
}
catch
{
Log ( $"Failed to download file '{contentId}'." );
throw ;
}
2023-04-13 09:33:10 +02:00
}
private void Log ( string msg )
{
2023-09-12 11:43:46 +02:00
tools . GetLog (). Log ( $"{GetName()}: {msg}" );
2023-04-13 09:33:10 +02:00
}
}
public class ContentId
{
public ContentId ( string id )
{
Id = id ;
}
public string Id { get ; }
2023-11-10 08:20:08 +01:00
public override bool Equals ( object? obj )
{
return obj is ContentId id && Id == id . Id ;
}
public override int GetHashCode ()
{
return HashCode . Combine ( Id );
}
2023-04-13 09:33:10 +02:00
}
}