Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
19 changes: 18 additions & 1 deletion src/Helpers/Constants.cs
Original file line number Diff line number Diff line change
Expand Up @@ -99,7 +99,7 @@ public class Constants
public static readonly decimal MAXIMUM_WITHDRAWAL_BTC_AMOUNT = 21_000_000;
public static readonly int TRANSACTION_CONFIRMATION_MINIMUM_BLOCKS;
public static int DEFAULT_CHANNEL_FEE_POLICY_TIMELOCK_DELTA_BLOCKS = 40;
public static long DEFAULT_CHANNEL_FEE_POLICY_BASE_FEE_MSAT = 0;
public static long DEFAULT_CHANNEL_FEE_POLICY_BASE_FEE_MSAT = 0;
public static long DEFAULT_CHANNEL_FEE_POLICY_FEE_RATE_PPM = 1500;
public static readonly long ANCHOR_CLOSINGS_MINIMUM_SATS;
public static readonly long MINIMUM_SWEEP_TRANSACTION_AMOUNT_SATS = 25_000_000; //25M sats
Expand Down Expand Up @@ -328,6 +328,13 @@ public class Constants
// Outbound ppm baseline for not-yet-categorized channels (safe mid default).
public static uint ROUTING_ENGINE_FEE_BASELINE_PPM_UNCATEGORIZED = 1500;

/// <summary>
/// Fraction of a channel's capacity advertised as its max_htlc_msat. LND itself uses ~0.99 for
/// channels it opens, so at the default this reconciles drifted channels without touching
/// untouched ones.
/// </summary>
public static double MAX_HTLC_CAPACITY_RATIO = 0.99;

public const string IsFrozenTag = "frozen";
public const string IsManuallyFrozenTag = "manually_frozen";

Expand Down Expand Up @@ -649,6 +656,16 @@ static Constants()
var feeBaselineUncategorized = Environment.GetEnvironmentVariable("ROUTING_ENGINE_FEE_BASELINE_PPM_UNCATEGORIZED");
if (feeBaselineUncategorized != null) ROUTING_ENGINE_FEE_BASELINE_PPM_UNCATEGORIZED = uint.Parse(feeBaselineUncategorized);

// Max HTLC
var maxHtlcCapacityRatio = Environment.GetEnvironmentVariable("MAX_HTLC_CAPACITY_RATIO");
if (maxHtlcCapacityRatio != null)
{
var parsedRatio = double.Parse(maxHtlcCapacityRatio, NumberStyles.AllowDecimalPoint | NumberStyles.AllowLeadingSign, CultureInfo.InvariantCulture);
// A ratio outside (0, 1] would resolve to 0 or above capacity, both of which LND rejects.
if (parsedRatio > 0 && parsedRatio <= 1) MAX_HTLC_CAPACITY_RATIO = parsedRatio;
else throw new ArgumentOutOfRangeException(nameof(MAX_HTLC_CAPACITY_RATIO), parsedRatio, "MAX_HTLC_CAPACITY_RATIO must be in (0, 1]");
}

// DB Initialization
ALICE_PUBKEY = Environment.GetEnvironmentVariable("ALICE_PUBKEY") ?? ALICE_PUBKEY;
ALICE_HOST = Environment.GetEnvironmentVariable("ALICE_HOST") ?? ALICE_HOST;
Expand Down
16 changes: 13 additions & 3 deletions src/Jobs/ChannelMonitorJob.cs
Original file line number Diff line number Diff line change
Expand Up @@ -87,6 +87,15 @@ public async Task Execute(IJobExecutionContext context)
// Recover Operations on channels
await RecoverGhostChannels(node1, node2, channel);
await RecoverChannelInConfirmationPendingStatus(node1);

try
{
await _lightningService.SyncChannelMaxHtlc(node1, channel);
}
catch (Exception e)
{
_logger.LogError(e, "Error while syncing max htlc for channel {ChanId} of node {NodeId}", channel.ChanId, node1.Id);
}
}
}
catch (Exception e)
Expand Down Expand Up @@ -124,7 +133,7 @@ private async Task RefreshExternalNodeData(Node managedNode, Node remoteNode, Li
return;
}

if (remoteNode.Name == nodeInfo.Alias) return;
if (remoteNode.Name == nodeInfo.Alias) return;
remoteNode.Name = nodeInfo.Alias;
var (updated, error) = _nodeRepository.Update(remoteNode);
if (!updated)
Expand All @@ -140,7 +149,7 @@ public async Task RecoverGhostChannels(Node source, Node destination, Channel ch
try
{
await using var dbContext = await _dbContextFactory.CreateDbContextAsync();

var channelPoint = channel.ChannelPoint.Split(":");
var fundingTx = channelPoint[0];
var outputIndex = Convert.ToUInt32(channelPoint[1]);
Expand All @@ -150,7 +159,8 @@ public async Task RecoverGhostChannels(Node source, Node destination, Channel ch

var parsedChannelPoint = new ChannelPoint
{
FundingTxidStr = fundingTx, FundingTxidBytes = ByteString.CopyFrom(Convert.FromHexString(fundingTx).Reverse().ToArray()),
FundingTxidStr = fundingTx,
FundingTxidBytes = ByteString.CopyFrom(Convert.FromHexString(fundingTx).Reverse().ToArray()),
OutputIndex = outputIndex
};

Expand Down
10 changes: 7 additions & 3 deletions src/Services/LightningClientService.cs
Original file line number Diff line number Diff line change
Expand Up @@ -48,8 +48,7 @@ public interface ILightningClientService
public void FundingStateStepVerify(Node node, PSBT finalizedPSBT, byte[] pendingChannelId, Lightning.LightningClient? client = null);
public void FundingStateStepFinalize(Node node, PSBT finalizedPSBT, byte[] pendingChannelId, Lightning.LightningClient? client = null);
public void FundingStateStepCancel(Node node, byte[] pendingChannelId, Lightning.LightningClient? client = null);

public Task<PolicyUpdateResponse?> SetChannelFeePolicy(Node node, NBitcoin.OutPoint chanPoint, long baseFeeMsat, uint feeRatePpm, uint timeLockDelta, int? inboundBaseFeeMsat, int? inboundFeeRatePpm, Lightning.LightningClient? client = null);
public Task<PolicyUpdateResponse?> SetChannelFeePolicy(Node node, NBitcoin.OutPoint chanPoint, long baseFeeMsat, uint feeRatePpm, uint timeLockDelta, int? inboundBaseFeeMsat, int? inboundFeeRatePpm, ulong? maxHtlcMsat = null, Lightning.LightningClient? client = null);
}

public class LightningClientService : ILightningClientService
Expand Down Expand Up @@ -453,7 +452,7 @@ public void FundingStateStepCancel(Node node, byte[] pendingChannelId, Lightning
}, new Metadata { { "macaroon", node.ChannelAdminMacaroon } });
}

public async Task<PolicyUpdateResponse?> SetChannelFeePolicy(Node node, NBitcoin.OutPoint chanPoint, long baseFeeMsat, uint feeRatePpm, uint timeLockDelta, int? inboundBaseFeeMsat, int? inboundFeeRatePpm, Lightning.LightningClient? client = null)
public async Task<PolicyUpdateResponse?> SetChannelFeePolicy(Node node, NBitcoin.OutPoint chanPoint, long baseFeeMsat, uint feeRatePpm, uint timeLockDelta, int? inboundBaseFeeMsat, int? inboundFeeRatePpm, ulong? maxHtlcMsat = null, Lightning.LightningClient? client = null)
{
client ??= GetLightningClient(node.Endpoint);

Expand All @@ -478,6 +477,11 @@ public void FundingStateStepCancel(Node node, byte[] pendingChannelId, Lightning
};
}

if (maxHtlcMsat.HasValue)
{
request.MaxHtlcMsat = maxHtlcMsat.Value;
}

return await client.UpdateChannelPolicyAsync(request, new Metadata { { "macaroon", node.ChannelAdminMacaroon } });
}
}
146 changes: 146 additions & 0 deletions src/Services/LightningService.cs
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,22 @@

namespace NodeGuard.Services
{
/// <summary>
/// Outcome of a <see cref="ILightningService.SyncChannelMaxHtlc"/> call. A failed write is an
/// exception rather than a value here: <see cref="Skipped"/> means we decided not to act.
/// </summary>
public enum MaxHtlcSyncResult
{
/// <summary>The channel already advertises the target max_htlc_msat — no RPC was made.</summary>
NoOp = 0,

/// <summary>The channel's max_htlc_msat was written to LND.</summary>
Updated = 1,

/// <summary>The target could not be resolved or the channel is not one we act on.</summary>
Skipped = 2,
}

/// <summary>
/// Service to interact with LND
/// </summary>
Expand Down Expand Up @@ -209,6 +225,15 @@ Task<Payment> SendPaymentV2Async(Node node, string paymentRequest, long amountSa
/// <param name="nodePubKey"></param>
/// <returns></returns>
public Task<(RoutingPolicy?, RoutingPolicy?)> GetChannelFeePolicy(ulong chanId, Node node);

/// <summary>
/// Reconciles the max_htlc_msat that <paramref name="node"/> advertises on
/// <paramref name="lndChannel"/> with <see cref="Constants.MAX_HTLC_CAPACITY_RATIO"/> of the
/// channel's capacity, writing to LND only when the advertised value differs.
/// </summary>
/// <param name="node">The managed node whose side of the channel is updated.</param>
/// <param name="lndChannel">The channel as reported by LND — the authority on capacity and chan id.</param>
public Task<MaxHtlcSyncResult> SyncChannelMaxHtlc(Node node, Lnrpc.Channel lndChannel);
}

public class LightningService : ILightningService
Expand Down Expand Up @@ -1894,5 +1919,126 @@ await _auditService.LogAsync(
return (managedNodePolicy, counterpartyNodePolicy);

}

public async Task<MaxHtlcSyncResult> SyncChannelMaxHtlc(Node node, Lnrpc.Channel lndChannel)
{
ArgumentNullException.ThrowIfNull(node);
ArgumentNullException.ThrowIfNull(lndChannel);

if (!node.IsManaged || string.IsNullOrWhiteSpace(node.ChannelAdminMacaroon))
{
_logger.LogWarning("Skipping max htlc sync for channel {ChanId}: node {NodeName} is not managed with channel admin access",
lndChannel.ChanId, node.Name);
return MaxHtlcSyncResult.Skipped;
}

if (!OutPoint.TryParse(lndChannel.ChannelPoint, out var outPoint))
{
_logger.LogWarning("Skipping max htlc sync for channel {ChanId} on {NodeName}: invalid chanPoint {ChanPoint}",
lndChannel.ChanId, node.Name, lndChannel.ChannelPoint);
return MaxHtlcSyncResult.Skipped;
}

RoutingPolicy? managedPolicy;
try
{
(managedPolicy, _) = await GetChannelFeePolicy(lndChannel.ChanId, node);
}
Comment thread
markettes marked this conversation as resolved.
catch (Exception e)
{
// A channel with no graph edge yet (freshly confirmed, or unannounced) throws here.
// The next monitor pass retries it.
_logger.LogWarning(e, "Skipping max htlc sync for channel {ChanId} on {NodeName}: current policy unavailable",
lndChannel.ChanId, node.Name);
return MaxHtlcSyncResult.Skipped;
}

if (managedPolicy == null)
{
_logger.LogWarning("Skipping max htlc sync for channel {ChanId} on {NodeName}: no policy for the managed side",
lndChannel.ChanId, node.Name);
return MaxHtlcSyncResult.Skipped;
}

var capacityMsat = (ulong)lndChannel.Capacity * 1_000;
var minHtlcMsat = (ulong)Math.Max(managedPolicy.MinHtlc, 0);

if (capacityMsat == 0 || minHtlcMsat > capacityMsat)
{
_logger.LogWarning("Skipping max htlc sync for channel {ChanId} on {NodeName}: no valid target between min_htlc {MinHtlcMsat} msat and capacity {CapacityMsat} msat",
lndChannel.ChanId, node.Name, minHtlcMsat, capacityMsat);
return MaxHtlcSyncResult.Skipped;
}
Comment thread
markettes marked this conversation as resolved.

var desiredMaxHtlcMsat = Math.Clamp(
(ulong)(capacityMsat * Constants.MAX_HTLC_CAPACITY_RATIO),
minHtlcMsat,
capacityMsat);

if (managedPolicy.MaxHtlcMsat == desiredMaxHtlcMsat)
{
_logger.LogDebug("Channel {ChanId} on {NodeName} already advertises max htlc {MaxHtlcMsat} msat",
lndChannel.ChanId, node.Name, desiredMaxHtlcMsat);
return MaxHtlcSyncResult.NoOp;
}
Comment thread
markettes marked this conversation as resolved.

// Only channels NodeGuard tracks are acted on, so the write is always auditable against a
// channel row.
var channel = await _channelRepository.GetByOutpoint(outPoint);
if (channel == null)
{
_logger.LogWarning("Skipping max htlc sync for channel {ChanId} on {NodeName}: no channel found for chanPoint {ChanPoint}",
lndChannel.ChanId, node.Name, lndChannel.ChannelPoint);
return MaxHtlcSyncResult.Skipped;
Comment thread
markettes marked this conversation as resolved.
}

// The fee fields are not being changed, but LND requires them to be echoed back in a policy update.
// Except the inbound fees, which are omitted to retain the current inbound policy.
var response = await _lightningClientService.SetChannelFeePolicy(
node,
outPoint,
managedPolicy.FeeBaseMsat,
(uint)Math.Clamp(managedPolicy.FeeRateMilliMsat, 0, uint.MaxValue),
managedPolicy.TimeLockDelta,
inboundBaseFeeMsat: null,
inboundFeeRatePpm: null,
maxHtlcMsat: desiredMaxHtlcMsat);

if (response?.FailedUpdates != null && response.FailedUpdates.Count > 0)
{
_logger.LogError("Failed to update max htlc for channel: {ChanPoint}", lndChannel.ChannelPoint);
throw new Exception($"Failed to update max htlc for channel: {lndChannel.ChannelPoint}");
}
Comment thread
markettes marked this conversation as resolved.

_logger.LogInformation("{NodeName} chan {ChanId}: set max htlc {PreviousMaxHtlcMsat}->{MaxHtlcMsat} msat (capacity {CapacityMsat} msat, ratio {Ratio})",
node.Name, lndChannel.ChanId, managedPolicy.MaxHtlcMsat, desiredMaxHtlcMsat, capacityMsat, Constants.MAX_HTLC_CAPACITY_RATIO);

try
{
await _auditService.LogSystemAsync(
AuditActionType.Update,
AuditEventType.Success,
AuditObjectType.Channel,
channel.Id.ToString(),
new
{
ChanPoint = lndChannel.ChannelPoint,
ChannelId = channel.Id,
lndChannel.ChanId,
NodeId = node.Id,
NodePubKey = node.PubKey,
PreviousMaxHtlcMsat = managedPolicy.MaxHtlcMsat,
MaxHtlcMsat = desiredMaxHtlcMsat,
CapacityMsat = capacityMsat,
CapacityRatio = Constants.MAX_HTLC_CAPACITY_RATIO
});
}
catch (Exception e)
{
_logger.LogError(e, "Error while saving max htlc audit log for chanPoint: {ChanPoint}", lndChannel.ChannelPoint);
}

return MaxHtlcSyncResult.Updated;
}
}
}
79 changes: 79 additions & 0 deletions test/NodeGuard.Tests/Jobs/ChannelMonitorJobTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,15 @@ private Mock<IDbContextFactory<ApplicationDbContext>> SetupDbContextFactory()
return dbContextFactory;
}

private Quartz.IJobExecutionContext BuildJobContext(int nodeId)
{
var jobDetail = new Mock<Quartz.IJobDetail>();
jobDetail.Setup(x => x.JobDataMap).Returns(new Quartz.JobDataMap { { "nodeId", nodeId.ToString() } });
var context = new Mock<Quartz.IJobExecutionContext>();
context.Setup(x => x.JobDetail).Returns(jobDetail.Object);
return context.Object;
}

[Fact]
public async Task RecoverGhostChannels_ChannelIsNotInitiatorButManaged()
{
Expand Down Expand Up @@ -227,6 +236,76 @@ public async Task RecoverGhostChannels_CreatesChannelNotInitiator()
context.Channels.Count().Should().Be(1);
}

[Fact]
public async Task Execute_SyncsMaxHtlcOfEveryChannel()
{
// Arrange
var logger = new Mock<ILogger<ChannelMonitorJob>>();
var dbContextFactory = SetupDbContextFactory();

var source = new Node() { Id = 3, Endpoint = "localhost", ChannelAdminMacaroon = "abc" };
// A managed peer we did not initiate with: ghost recovery and alias refresh both bail out early,
// leaving the max htlc sync as the only work Execute does per channel.
var remote = new Node() { Id = 9, PubKey = "peer", Endpoint = "localhost" };
var channel1 = new Lnrpc.Channel() { ChanId = 1, Capacity = 1000, RemotePubkey = remote.PubKey, Initiator = false };
var channel2 = new Lnrpc.Channel() { ChanId = 2, Capacity = 2000, RemotePubkey = remote.PubKey, Initiator = false };

var nodeRepository = new Mock<NodeGuard.Data.Repositories.Interfaces.INodeRepository>();
nodeRepository.Setup(x => x.GetById(source.Id)).ReturnsAsync(source);
nodeRepository.Setup(x => x.GetOrCreateByPubKey(remote.PubKey, It.IsAny<ILightningService>())).ReturnsAsync(remote);

var lightningClientService = new Mock<ILightningClientService>();
lightningClientService.Setup(x => x.ListChannels(source, It.IsAny<Lightning.LightningClient>()))
.ReturnsAsync(new ListChannelsResponse { Channels = { channel1, channel2 } });

var lightningService = new Mock<ILightningService>();
lightningService.Setup(x => x.SyncChannelMaxHtlc(source, It.IsAny<Lnrpc.Channel>())).ReturnsAsync(MaxHtlcSyncResult.Updated);

var channelMonitorJob = new ChannelMonitorJob(logger.Object, dbContextFactory.Object, nodeRepository.Object, lightningService.Object, lightningClientService.Object);

// Act
var act = () => channelMonitorJob.Execute(BuildJobContext(source.Id));

// Assert
await act.Should().NotThrowAsync();
lightningService.Verify(x => x.SyncChannelMaxHtlc(source, channel1), Times.Once);
lightningService.Verify(x => x.SyncChannelMaxHtlc(source, channel2), Times.Once);
}

[Fact]
public async Task Execute_MaxHtlcSyncThrows()
{
// Arrange
var logger = new Mock<ILogger<ChannelMonitorJob>>();
var dbContextFactory = SetupDbContextFactory();

var source = new Node() { Id = 3, Endpoint = "localhost", ChannelAdminMacaroon = "abc" };
var remote = new Node() { Id = 9, PubKey = "peer", Endpoint = "localhost" };
var channel1 = new Lnrpc.Channel() { ChanId = 1, Capacity = 1000, RemotePubkey = remote.PubKey, Initiator = false };
var channel2 = new Lnrpc.Channel() { ChanId = 2, Capacity = 2000, RemotePubkey = remote.PubKey, Initiator = false };

var nodeRepository = new Mock<NodeGuard.Data.Repositories.Interfaces.INodeRepository>();
nodeRepository.Setup(x => x.GetById(source.Id)).ReturnsAsync(source);
nodeRepository.Setup(x => x.GetOrCreateByPubKey(remote.PubKey, It.IsAny<ILightningService>())).ReturnsAsync(remote);

var lightningClientService = new Mock<ILightningClientService>();
lightningClientService.Setup(x => x.ListChannels(source, It.IsAny<Lightning.LightningClient>()))
.ReturnsAsync(new ListChannelsResponse { Channels = { channel1, channel2 } });

var lightningService = new Mock<ILightningService>();
lightningService.Setup(x => x.SyncChannelMaxHtlc(source, channel1)).ThrowsAsync(new Exception("policy update rejected"));
lightningService.Setup(x => x.SyncChannelMaxHtlc(source, channel2)).ReturnsAsync(MaxHtlcSyncResult.Updated);

var channelMonitorJob = new ChannelMonitorJob(logger.Object, dbContextFactory.Object, nodeRepository.Object, lightningService.Object, lightningClientService.Object);

// Act
var act = () => channelMonitorJob.Execute(BuildJobContext(source.Id));

// Assert - a failed policy write is contained, so the run finishes and the next channel is synced
await act.Should().NotThrowAsync();
lightningService.Verify(x => x.SyncChannelMaxHtlc(source, channel2), Times.Once);
}

[Fact]
public async Task RecoverChannelInConfirmationPendingStatus_RequestWithDifferentSource()
{
Expand Down
4 changes: 2 additions & 2 deletions test/NodeGuard.Tests/Services/LightningClientServiceTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -99,7 +99,7 @@ public async Task SetChannelFeePolicy_BuildsPolicyUpdateRequestWithInboundFee()
timeLockDelta: 40,
inboundBaseFeeMsat: -100,
inboundFeeRatePpm: -25,
lightningClient.Object);
client: lightningClient.Object);

// Assert
response.Should().NotBeNull();
Expand Down Expand Up @@ -152,7 +152,7 @@ await lightningClientService.SetChannelFeePolicy(
timeLockDelta: 40,
inboundBaseFeeMsat: null,
inboundFeeRatePpm: null,
lightningClient.Object);
client: lightningClient.Object);

// Assert
capturedRequest.Should().NotBeNull();
Expand Down
Loading
Loading