Schreiben Sie Ihren ersten Trigger Node
Dieses Tutorial hilft Ihnen dabei, Ihren ersten Trigger Node für OctoMesh zu schreiben.
Der Kürze halber gehen wir davon aus, dass Sie Ihre Entwicklungsumgebung bereits eingerichtet haben und über Grundkenntnisse in C# und .NET Core verfügen.
Wir gehen außerdem davon aus, dass Sie einen Edge oder Mesh Adapter mithilfe des Tutorials Write your first Edge Adapter oder Write your first Mesh Adapter eingerichtet haben.
Eine Referenzimplementierung finden Sie im folgenden Repository
Überblick
Trigger Nodes werden verwendet, um den Transformationsteil im OctoMesh-System auszulösen. Trigger basieren auf Ereignissen (intern oder extern). Eine Pipeline-Definition kann mehrere Trigger enthalten, benötigt aber mindestens einen Trigger, um die Pipeline-Ausführung zu starten.
Beispiele für Trigger sind:
- FromExecutePipelineCommand: Ein Trigger, der eine Pipeline-Ausführung durch einen Befehl aus der AdminUI oder CLI startet.
- FromHttpRequest: Ein Trigger, der eine Pipeline-Ausführung durch eine HTTP-Anfrage startet.
- FromPipelineTriggerEvent: Ein Trigger, der eine Pipeline-Ausführung durch ein Ereignis (vom BotService) startet.
Trigger starten die Ausführung des Transformationsteils der Pipeline. Sie können dem Transformationsteil Eingabedaten übergeben. Ebenso kann die Ausgabe des Transformationsteils an den Trigger zurückgegeben und weiter an ein externes System übergeben werden.
In diesem Tutorial erstellen wir einen Trigger, der eine Pipeline-Ausführung auf Basis einer TCP-Anfrage startet. Dazu müssen wir die folgenden Schritte ausführen:
1. Create a new node configuration
Erstellen Sie eine neue Klasse, die die Klasse TriggerNodeConfiguration aus dem OctoMesh SDK erweitert.
using Meshmakers.Octo.Sdk.Common;
[NodeName("DemoTrigger", 1)]
public record DemoTriggerNodeConfiguration : TriggerNodeConfiguration
{
/// <summary>
/// Port the sample TCP listener listens to
/// </summary>
public required ushort Port { get; set; } = 8000;
}
Diese Node-Konfigurationsklasse wird verwendet, um den Trigger Node zu konfigurieren.
Das Attribut NodeName wird verwendet, um den Namen des Nodes und die Version des Nodes zu definieren.
Es ist wichtig, die Version des Nodes zu definieren, denn die Version wird verwendet, um zwischen verschiedenen Versionen desselben Nodes zu unterscheiden.
Wir definieren eine Eigenschaft Port, die verwendet wird, um den Port zu konfigurieren, auf dem der TCP-Listener lauscht. Der Standardwert ist 8000.
2. Create the trigger node
Erstellen Sie eine neue Klasse, die die Schnittstelle ITriggerNode aus dem OctoMesh SDK implementiert.
[NodeConfiguration(typeof(DemoTriggerNodeConfiguration))]
// ReSharper disable once ClassNeverInstantiated.Global
public class DemoTriggerNode : ITriggerPipelineNode
{
public Task StartAsync(ITriggerContext context)
{
return Task.CompletedTask;
}
public Task StopAsync(ITriggerContext context)
{
return Task.CompletedTask;
}
}
Das Attribut NodeConfiguration wird verwendet, um die Konfigurationsklasse für den Node zu definieren.
Die Klasse DemoTriggerNode implementiert die Schnittstelle ITriggerPipelineNode.
Die Methode StartAsync wird aufgerufen, wenn der Trigger Node gestartet wird. Die Methode StopAsync wird aufgerufen, wenn der Trigger Node stoppt.
Trigger Nodes werden vom OctoMesh-Adapter gestartet und gestoppt, wenn der Adapter gestartet und gestoppt wird.
Dies geschieht durch die IAdapterService-Implementierung des Adapters über den Aufruf StartTriggerPipelineNodesAsync der Schnittstelle IPipelineRegistryService.
Der Start und Stopp erfolgt aus diesem Grund auch, wenn eine neue Konfiguration auf den Adapter ausgerollt wird.
3. Implement the trigger logic
Die Trigger-Logik muss einen TCP-Listener öffnen und auf eingehende Verbindungen warten. Wenn eine Verbindung hergestellt wird, muss der Trigger Node die Pipeline-Ausführung starten. Die Eingabedaten werden an die Pipeline-Ausführung übergeben, und die Ausgabe wird als Ergebnis der TCP-Anfrage zurückgegeben.
Öffnen wir den TCP-Listener und warten wir in der Methode StartAsync auf eingehende Verbindungen. In der Methode StopAsync stoppen wir den TCP-Listener.
```sharp
[NodeConfiguration(typeof(DemoTriggerNodeConfiguration))]
// ReSharper disable once ClassNeverInstantiated.Global
public class DemoTriggerNode( /* you can inject services here */) : ITriggerPipelineNode
{
private TcpListener? _tcpListener;
private CancellationTokenSource? _cts;
public Task StartAsync(ITriggerContext context)
{
var c = context.NodeContext.GetNodeConfiguration<DemoTriggerNodeConfiguration>();
// Listen to TCP port for incoming messages at all interfaces
_cts = new CancellationTokenSource();
_tcpListener = new TcpListener(IPAddress.Any, c.Port);
_tcpListener.Start();
context.NodeContext.Info("TCP listener started and waiting for connections...");
Task.Run(async () =>
{
// Wait for incoming connections
while (!_cts.Token.IsCancellationRequested)
{
try
{
var tcpClient = await _tcpListener.AcceptTcpClientAsync();
_ = ProcessClientAsync(tcpClient, context, _cts.Token); // Execute in a separate task
}
catch (ObjectDisposedException)
{
// The Listener was stopped
break;
}
}
});
return Task.CompletedTask;
}
public Task StopAsync(ITriggerContext context)
{
try
{
_cts?.Cancel();
_tcpListener?.Stop();
context.NodeContext.Info("TCP listener stopped.");
}
catch (Exception ex)
{
context.NodeContext.Error($"Error stopping: {ex.Message}");
throw DemoPipelineExecutionException.StopFailed(ex);
}
return Task.CompletedTask;
}
private async Task ProcessClientAsync(TcpClient client, ITriggerContext context,
CancellationToken cancellationToken)
{
return Task.CompletedTask;
}
}
Wir müssen die eingehende Nachricht als JSON parsen und die Pipeline-Ausführung in der Methode ProcessClientAsync starten. Wir führen etwas Ausnahmebehandlung für das JSON-Parsing und die Pipeline-Ausführung durch.
using Meshmakers.Octo.Sdk.Common.Services;
using Newtonsoft.Json;
namespace Meshmakers.Octo.Communication.MeshAdapter.Demo;
public class DemoPipelineExecutionException : PipelineExecutionException
{
public DemoPipelineExecutionException()
{
}
public DemoPipelineExecutionException(string message) : base(message)
{
}
public DemoPipelineExecutionException(string message, Exception inner) : base(message, inner)
{
}
public static Exception StopFailed(Exception exception)
{
return new DemoPipelineExecutionException("Stop failed", exception);
}
public static Exception PipelineExecutionFailed(Exception exception)
{
return new DemoPipelineExecutionException("Execution of pipeline failed", exception);
}
public static Exception MessageDeserializationFailed(JsonReaderException jsonReaderException)
{
return new DemoPipelineExecutionException("Message deserialization failed", jsonReaderException);
}
}
Der nächste Schritt ist die Implementierung der Methode ProcessClientAsync. Diese Methode liest eingehende Nachrichten vom Client, startet die Pipeline-Ausführung und sendet die Ausgabe zurück an den Client.
Wir müssen die eingehende Nachricht als JSON parsen und die Pipeline-Ausführung mit der Nachricht als Eingabe starten. Wir führen etwas Ausnahmebehandlung für das JSON-Parsing und die Pipeline-Ausführung durch.
private async Task ProcessClientAsync(TcpClient client, ITriggerContext context, CancellationToken cancellationToken)
{
try
{
// Read incoming messages from the client
await using var networkStream = client.GetStream();
var buffer = new byte[1024];
int bytesRead;
var messageBuilder = new StringBuilder();
while ((bytesRead = await networkStream.ReadAsync(buffer, 0, buffer.Length, cancellationToken)) > 0)
{
context.NodeContext.Info("Client connected.");
// Process the incoming message
messageBuilder.Append(Encoding.UTF8.GetString(buffer, 0, bytesRead));
if (networkStream.DataAvailable)
{
continue;
}
var message = messageBuilder.ToString();
context.NodeContext.Info($"Received message: {message}");
// Parse the incoming message as JSON and start the pipeline execution with the message as input
var input = JToken.Parse(message);
var output = await context.ExecuteAsync(
new ExecutePipelineOptions(DateTime.UtcNow)
{
ExternalReceivedDateTime = DateTime.UtcNow
},
input);
// Serialize the output as JSON and send it back to the client
var outputString = JsonConvert.SerializeObject(output);
byte[] utf8Bytes = Encoding.UTF8.GetBytes(outputString + Environment.NewLine);
await networkStream.WriteAsync(utf8Bytes, cancellationToken);
await networkStream.FlushAsync(cancellationToken);
messageBuilder.Clear();
}
}
catch (JsonReaderException ex)
{
context.NodeContext.Error($"Error parsing input: {ex.Message}");
throw DemoPipelineExecutionException.MessageDeserializationFailed(ex);
}
catch (Exception ex)
{
context.NodeContext.Error($"Error processing message: {ex.Message}");
throw DemoPipelineExecutionException.PipelineExecutionFailed(ex);
}
finally
{
client.Close();
}
}
4. Create pipeline definition with the trigger node
Der letzte Schritt besteht darin, eine Pipeline-Definition zu erstellen, die den Trigger Node enthält. Die Pipeline-Definition ist eine YAML-Datei, die die Pipeline-Struktur definiert.
Wir erstellen eine neue Pipeline in der AdminUI und fügen den Trigger Node zur Pipeline-Definition hinzu.
triggers:
- type: DemoTrigger@1
port: 8000
transformations:
- type: Demo@1
description: Simulates data
myMessage: Hello, Mars!
Starten Sie den Adapter und rollen Sie die Pipeline-Definition auf den Adapter aus. Der Trigger Node beginnt, auf eingehende TCP-Verbindungen auf Port 8000 zu lauschen.
Um den Trigger Node zu testen, können Sie einen TCP-Client wie netcat oder telnet verwenden, um eine Nachricht an den Trigger Node zu senden.
echo '{"message": "Hello, Earth!"}' | nc localhost 8000
Sie können auch das vorbereitete PowerShell-Skript verwenden, um eine Nachricht an den Trigger Node zu senden.
Wenn der Trigger Node eine Nachricht empfängt, startet er die Pipeline-Ausführung mit der Nachricht als Eingabe. Die Ausgabe der Pipeline-Ausführung wird zurück an den Client gesendet. Die Ausgabe können Sie in der Konsole sehen, in der der Adapter läuft.
Hello, Mars!