diff --git a/Kromer/Controllers/V1/SubscriptionsController.cs b/Kromer/Controllers/V1/SubscriptionsController.cs
new file mode 100644
index 0000000..78949d8
--- /dev/null
+++ b/Kromer/Controllers/V1/SubscriptionsController.cs
@@ -0,0 +1,127 @@
+using Kromer.Models.Api.V1;
+using Kromer.Models.Api.V1.Subscriptions;
+using Kromer.Models.Dto;
+using Kromer.Models.Exceptions;
+using Kromer.Repositories;
+using Microsoft.AspNetCore.Mvc;
+
+namespace Kromer.Controllers.V1;
+
+[Route("api/v1/subscriptions")]
+[ApiController]
+public class SubscriptionsController(SubscriptionRepository subscriptionRepository) : ControllerBase
+{
+ ///
+ /// Creates a subscription contract for the supplied name or metaname.
+ ///
+ /// The contract details and private key of the current name owner.
+ /// The identifier of the created subscription contract.
+ /// Thrown when the request is invalid, authentication fails, or the name does not exist.
+ [HttpPost("")]
+ public async Task>> CreateSubscription(
+ [FromBody] CreateSubscriptionRequest? request)
+ {
+ return new Result(
+ await subscriptionRepository.CreateContractAsync(request));
+ }
+
+ ///
+ /// Cancels a subscription contract and all active wallet subscriptions attached to it.
+ ///
+ /// The identifier of the subscription contract to cancel.
+ /// The private key of the current name owner.
+ /// The cancelled subscription contract.
+ /// Thrown when the contract is not found, the private key is invalid, or the caller does not own the name.
+ [HttpDelete("{id:int}")]
+ public async Task>> DeleteSubscription(int id,
+ [FromBody] PrivateKeyRequest? request)
+ {
+ return new Result(
+ await subscriptionRepository.CancelContractAsync(id, request?.PrivateKey));
+ }
+
+ ///
+ /// Closes a subscription contract to new subscribers while keeping existing subscriptions active.
+ ///
+ /// The identifier of the subscription contract to close.
+ /// The private key of the current name owner.
+ /// The closed subscription contract.
+ /// Thrown when the contract is not found, the private key is invalid, or the caller does not own the name.
+ [HttpPost("{id:int}/close")]
+ public async Task>> CloseSubscription(int id,
+ [FromBody] PrivateKeyRequest? request)
+ {
+ return new Result(
+ await subscriptionRepository.CloseContractAsync(id, request?.PrivateKey));
+ }
+
+ ///
+ /// Retrieves a subscription contract by its identifier.
+ ///
+ /// The identifier of the subscription contract to retrieve.
+ /// Optional wallet addresses used to include caller-specific subscription state.
+ /// The subscription contract details.
+ /// Thrown when the contract is not found or the address is invalid.
+ [HttpGet("{id:int}")]
+ public async Task>> GetSubscription(int id,
+ [FromQuery(Name = "address")] List? addresses = null)
+ {
+ return new Result(
+ await subscriptionRepository.GetContractAsync(id, addresses));
+ }
+
+ ///
+ /// Subscribes the authenticated wallet to a subscription contract.
+ ///
+ /// The identifier of the subscription contract to subscribe to.
+ /// The private key of the subscribing wallet.
+ /// The next payment date for the wallet subscription.
+ /// Thrown when authentication fails or the wallet cannot subscribe to the contract.
+ [HttpPost("{id:int}/subscribe")]
+ public async Task>> Subscribe(int id, [FromBody] PrivateKeyRequest? request)
+ {
+ return new Result(
+ await subscriptionRepository.SubscribeAsync(id, request?.PrivateKey));
+ }
+
+ ///
+ /// Unsubscribes the authenticated wallet from a subscription contract.
+ ///
+ /// The identifier of the subscription contract to unsubscribe from.
+ /// The private key of the subscribed wallet.
+ /// An empty result when the unsubscribe request has been processed.
+ /// Thrown when authentication fails, the contract is not found, or the subscription cannot be unsubscribed.
+ [HttpPost("{id:int}/unsubscribe")]
+ public async Task>> Unsubscribe(int id, [FromBody] PrivateKeyRequest? request)
+ {
+ await subscriptionRepository.UnsubscribeAsync(id, request?.PrivateKey);
+ return new Result(new { });
+ }
+
+ ///
+ /// Lists subscription contracts related to a wallet address or name.
+ ///
+ /// The wallet addresses used to find owned and subscribed contracts.
+ /// The names or metanames used to filter contracts by receiver.
+ /// Excludes contracts owned by the supplied addresses when set.
+ /// Only includes contracts owned by the supplied addresses when set.
+ /// Only includes active wallet subscriptions that can currently be unsubscribed.
+ /// The maximum number of contracts to return. The value is between 1 and 1000!
+ /// The number of contracts to skip before returning results.
+ /// A paginated list of subscription contracts.
+ /// Thrown when neither address nor name is supplied, or when supplied filters are invalid.
+ [HttpGet("")]
+ public async Task>> ListSubscriptions(
+ [FromQuery(Name = "address")] List? addresses = null,
+ [FromQuery(Name = "name")] List? names = null,
+ [FromQuery(Name = "exclude_owned")] bool excludeOwned = false,
+ [FromQuery(Name = "only_owned")] bool onlyOwned = false,
+ [FromQuery(Name = "only_unsubscribable")] bool onlyUnsubscribable = true,
+ [FromQuery] int limit = 50,
+ [FromQuery] int offset = 0)
+ {
+ return new Result(
+ await subscriptionRepository.ListContractsAsync(addresses, names, excludeOwned, onlyOwned,
+ onlyUnsubscribable, limit, offset));
+ }
+}
diff --git a/Kromer/Data/KromerContext.cs b/Kromer/Data/KromerContext.cs
index 9682c49..dba6eee 100644
--- a/Kromer/Data/KromerContext.cs
+++ b/Kromer/Data/KromerContext.cs
@@ -19,16 +19,21 @@ public KromerContext(DbContextOptions options)
public virtual DbSet Players { get; set; }
+ public virtual DbSet SubscriptionContracts { get; set; }
+
public virtual DbSet Transactions { get; set; }
public virtual DbSet Wallets { get; set; }
+ public virtual DbSet WalletSubscriptions { get; set; }
+
protected override void OnModelCreating(ModelBuilder modelBuilder)
{
modelBuilder
.HasPostgresEnum("transaction_type",
new[] { "mined", "unknown", "name_purchase", "name_a_record", "name_transfer", "transfer" })
- .HasPostgresEnum("transaction_type_enum", new[] { "unknown", "mined", "name_purchase", "transfer" });
+ .HasPostgresEnum("transaction_type_enum", new[] { "unknown", "mined", "name_purchase", "transfer" })
+ .HasPostgresEnum("subscription_status", new[] { "active", "closed", "cancelled" });
modelBuilder.Entity(entity =>
{
@@ -114,6 +119,41 @@ protected override void OnModelCreating(ModelBuilder modelBuilder)
.HasColumnType("public.transaction_type");
});
+ modelBuilder.Entity(entity =>
+ {
+ entity.HasKey(e => e.Id).HasName("subscription_contracts_pkey");
+
+ entity.ToTable("subscription_contracts");
+
+ entity.HasIndex(e => e.BaseName, "idx_subscription_contract_base_name");
+ entity.HasIndex(e => e.Receiver, "idx_subscription_contract_receiver");
+ entity.HasIndex(e => e.Status, "idx_subscription_contract_status");
+
+ entity.Property(e => e.Id).HasColumnName("id");
+ entity.Property(e => e.Receiver)
+ .HasMaxLength(98)
+ .HasColumnName("receiver");
+ entity.Property(e => e.BaseName)
+ .HasMaxLength(64)
+ .HasColumnName("base_name");
+ entity.Property(e => e.Price)
+ .HasPrecision(16, 5)
+ .HasColumnName("price");
+ entity.Property(e => e.PeriodMinutes).HasColumnName("period_minutes");
+ entity.Property(e => e.Description)
+ .HasMaxLength(255)
+ .HasColumnName("description");
+ entity.Property(e => e.MaxSubscribers).HasColumnName("max_subscribers");
+ entity.Property(e => e.AllowedSubscriberAddresses)
+ .HasColumnType("text[]")
+ .HasColumnName("allowed_subscriber_addresses");
+ entity.Property(e => e.Status)
+ .HasColumnName("status")
+ .HasColumnType("public.subscription_status");
+ entity.Property(e => e.CreatedAt).HasColumnName("created_at");
+ entity.Property(e => e.CancelledAt).HasColumnName("cancelled_at");
+ });
+
modelBuilder.Entity(entity =>
{
entity.HasKey(e => e.Id).HasName("wallets_pkey");
@@ -145,8 +185,47 @@ protected override void OnModelCreating(ModelBuilder modelBuilder)
.HasColumnName("total_out");
});
+ modelBuilder.Entity(entity =>
+ {
+ entity.HasKey(e => e.Id).HasName("wallet_subscriptions_pkey");
+
+ entity.ToTable("wallet_subscriptions");
+
+ entity.HasIndex(e => e.ContractId, "idx_wallet_subscription_contract_id");
+ entity.HasIndex(e => e.WalletAddress, "idx_wallet_subscription_wallet_address");
+ entity.HasIndex(e => e.NextPayment, "idx_wallet_subscription_next_payment");
+ entity.HasIndex(e => new { e.ContractId, e.WalletAddress }, "unique_active_wallet_subscription")
+ .IsUnique()
+ .HasFilter("status = 'active'");
+
+ entity.Property(e => e.Id).HasColumnName("id");
+ entity.Property(e => e.ContractId).HasColumnName("contract_id");
+ entity.Property(e => e.WalletAddress)
+ .HasMaxLength(10)
+ .IsFixedLength()
+ .HasColumnName("wallet_address");
+ entity.Property(e => e.NextPayment).HasColumnName("next_payment");
+ entity.Property(e => e.Status)
+ .HasColumnName("status")
+ .HasColumnType("public.subscription_status");
+ entity.Property(e => e.CancellationReason)
+ .HasMaxLength(64)
+ .HasColumnName("cancellation_reason");
+ entity.Property(e => e.CanUnsubscribe)
+ .HasDefaultValue(true)
+ .HasColumnName("can_unsubscribe");
+ entity.Property(e => e.CreatedAt).HasColumnName("created_at");
+ entity.Property(e => e.CancelledAt).HasColumnName("cancelled_at");
+
+ entity.HasOne(e => e.Contract)
+ .WithMany(e => e.WalletSubscriptions)
+ .HasForeignKey(e => e.ContractId)
+ .HasConstraintName("wallet_subscriptions_contract_id_fkey")
+ .OnDelete(DeleteBehavior.Cascade);
+ });
+
OnModelCreatingPartial(modelBuilder);
}
partial void OnModelCreatingPartial(ModelBuilder modelBuilder);
-}
\ No newline at end of file
+}
diff --git a/Kromer/Migrations/20260427000000_AddSubscriptions.cs b/Kromer/Migrations/20260427000000_AddSubscriptions.cs
new file mode 100644
index 0000000..13cc681
--- /dev/null
+++ b/Kromer/Migrations/20260427000000_AddSubscriptions.cs
@@ -0,0 +1,130 @@
+using Kromer.Data;
+using Kromer.Models.Entities;
+using Microsoft.EntityFrameworkCore.Infrastructure;
+using Microsoft.EntityFrameworkCore.Migrations;
+using Npgsql.EntityFrameworkCore.PostgreSQL.Metadata;
+
+#nullable disable
+
+namespace Kromer.Migrations;
+
+[DbContext(typeof(KromerContext))]
+[Migration("20260427000000_AddSubscriptions")]
+public partial class AddSubscriptions : Migration
+{
+ protected override void Up(MigrationBuilder migrationBuilder)
+ {
+ migrationBuilder.Sql("""
+ DO $$
+ BEGIN
+ IF NOT EXISTS (SELECT 1 FROM pg_type WHERE typname = 'subscription_status') THEN
+ CREATE TYPE public.subscription_status AS ENUM ('active', 'closed', 'cancelled');
+ ELSIF NOT EXISTS (
+ SELECT 1
+ FROM pg_enum e
+ JOIN pg_type t ON t.oid = e.enumtypid
+ WHERE t.typname = 'subscription_status' AND e.enumlabel = 'closed'
+ ) THEN
+ ALTER TYPE public.subscription_status ADD VALUE 'closed' BEFORE 'cancelled';
+ END IF;
+ END
+ $$;
+ """);
+
+ migrationBuilder.CreateTable(
+ name: "subscription_contracts",
+ columns: table => new
+ {
+ id = table.Column(type: "integer", nullable: false)
+ .Annotation("Npgsql:ValueGenerationStrategy", NpgsqlValueGenerationStrategy.IdentityByDefaultColumn),
+ receiver = table.Column(type: "character varying(98)", maxLength: 98, nullable: false),
+ base_name = table.Column(type: "character varying(64)", maxLength: 64, nullable: false),
+ price = table.Column(type: "numeric(16,5)", precision: 16, scale: 5, nullable: false),
+ period_minutes = table.Column(type: "integer", nullable: false),
+ description = table.Column(type: "character varying(255)", maxLength: 255, nullable: false),
+ max_subscribers = table.Column(type: "integer", nullable: true),
+ allowed_subscriber_addresses = table.Column(type: "text[]", nullable: true),
+ status = table.Column(type: "public.subscription_status", nullable: false),
+ created_at = table.Column(type: "timestamp with time zone", nullable: false),
+ cancelled_at = table.Column(type: "timestamp with time zone", nullable: true)
+ },
+ constraints: table =>
+ {
+ table.PrimaryKey("subscription_contracts_pkey", x => x.id);
+ });
+
+ migrationBuilder.CreateTable(
+ name: "wallet_subscriptions",
+ columns: table => new
+ {
+ id = table.Column(type: "integer", nullable: false)
+ .Annotation("Npgsql:ValueGenerationStrategy", NpgsqlValueGenerationStrategy.IdentityByDefaultColumn),
+ contract_id = table.Column(type: "integer", nullable: false),
+ wallet_address = table.Column(type: "character(10)", fixedLength: true, maxLength: 10, nullable: false),
+ next_payment = table.Column(type: "timestamp with time zone", nullable: false),
+ status = table.Column(type: "public.subscription_status", nullable: false),
+ cancellation_reason = table.Column(type: "character varying(64)", maxLength: 64, nullable: true),
+ can_unsubscribe = table.Column(type: "boolean", nullable: false, defaultValue: true),
+ created_at = table.Column(type: "timestamp with time zone", nullable: false),
+ cancelled_at = table.Column(type: "timestamp with time zone", nullable: true)
+ },
+ constraints: table =>
+ {
+ table.PrimaryKey("wallet_subscriptions_pkey", x => x.id);
+ table.ForeignKey(
+ name: "wallet_subscriptions_contract_id_fkey",
+ column: x => x.contract_id,
+ principalTable: "subscription_contracts",
+ principalColumn: "id",
+ onDelete: ReferentialAction.Cascade);
+ });
+
+ migrationBuilder.CreateIndex(
+ name: "idx_subscription_contract_base_name",
+ table: "subscription_contracts",
+ column: "base_name");
+
+ migrationBuilder.CreateIndex(
+ name: "idx_subscription_contract_receiver",
+ table: "subscription_contracts",
+ column: "receiver");
+
+ migrationBuilder.CreateIndex(
+ name: "idx_subscription_contract_status",
+ table: "subscription_contracts",
+ column: "status");
+
+ migrationBuilder.CreateIndex(
+ name: "idx_wallet_subscription_contract_id",
+ table: "wallet_subscriptions",
+ column: "contract_id");
+
+ migrationBuilder.CreateIndex(
+ name: "idx_wallet_subscription_next_payment",
+ table: "wallet_subscriptions",
+ column: "next_payment");
+
+ migrationBuilder.CreateIndex(
+ name: "idx_wallet_subscription_wallet_address",
+ table: "wallet_subscriptions",
+ column: "wallet_address");
+
+ migrationBuilder.CreateIndex(
+ name: "unique_active_wallet_subscription",
+ table: "wallet_subscriptions",
+ columns: new[] { "contract_id", "wallet_address" },
+ unique: true,
+ filter: "status = 'active'");
+ }
+
+ protected override void Down(MigrationBuilder migrationBuilder)
+ {
+ migrationBuilder.DropTable(
+ name: "wallet_subscriptions");
+
+ migrationBuilder.DropTable(
+ name: "subscription_contracts");
+
+ migrationBuilder.Sql("DROP TYPE IF EXISTS public.subscription_status;");
+ }
+}
diff --git a/Kromer/Models/Api/V1/Subscriptions/CreateSubscriptionRequest.cs b/Kromer/Models/Api/V1/Subscriptions/CreateSubscriptionRequest.cs
new file mode 100644
index 0000000..f1c859a
--- /dev/null
+++ b/Kromer/Models/Api/V1/Subscriptions/CreateSubscriptionRequest.cs
@@ -0,0 +1,21 @@
+using System.Text.Json.Serialization;
+
+namespace Kromer.Models.Api.V1.Subscriptions;
+
+public class CreateSubscriptionRequest
+{
+ [JsonPropertyName("privatekey")]
+ public required string PrivateKey { get; set; }
+
+ public required string Name { get; set; }
+
+ public decimal Price { get; set; }
+
+ public int Period { get; set; }
+
+ public required string Description { get; set; }
+
+ public int? MaxSubscribers { get; set; }
+
+ public List? AllowedSubscribers { get; set; }
+}
diff --git a/Kromer/Models/Api/V1/Subscriptions/CreateSubscriptionResponse.cs b/Kromer/Models/Api/V1/Subscriptions/CreateSubscriptionResponse.cs
new file mode 100644
index 0000000..82a0180
--- /dev/null
+++ b/Kromer/Models/Api/V1/Subscriptions/CreateSubscriptionResponse.cs
@@ -0,0 +1,6 @@
+namespace Kromer.Models.Api.V1.Subscriptions;
+
+public class CreateSubscriptionResponse
+{
+ public int Id { get; set; }
+}
diff --git a/Kromer/Models/Api/V1/Subscriptions/PrivateKeyRequest.cs b/Kromer/Models/Api/V1/Subscriptions/PrivateKeyRequest.cs
new file mode 100644
index 0000000..d37afbf
--- /dev/null
+++ b/Kromer/Models/Api/V1/Subscriptions/PrivateKeyRequest.cs
@@ -0,0 +1,9 @@
+using System.Text.Json.Serialization;
+
+namespace Kromer.Models.Api.V1.Subscriptions;
+
+public class PrivateKeyRequest
+{
+ [JsonPropertyName("privatekey")]
+ public required string PrivateKey { get; set; }
+}
diff --git a/Kromer/Models/Api/V1/Subscriptions/SubscribeResponse.cs b/Kromer/Models/Api/V1/Subscriptions/SubscribeResponse.cs
new file mode 100644
index 0000000..6bd30fe
--- /dev/null
+++ b/Kromer/Models/Api/V1/Subscriptions/SubscribeResponse.cs
@@ -0,0 +1,6 @@
+namespace Kromer.Models.Api.V1.Subscriptions;
+
+public class SubscribeResponse
+{
+ public DateTime NextPayment { get; set; }
+}
diff --git a/Kromer/Models/Api/V1/Subscriptions/SubscriptionListResponse.cs b/Kromer/Models/Api/V1/Subscriptions/SubscriptionListResponse.cs
new file mode 100644
index 0000000..4509efa
--- /dev/null
+++ b/Kromer/Models/Api/V1/Subscriptions/SubscriptionListResponse.cs
@@ -0,0 +1,12 @@
+using Kromer.Models.Dto;
+
+namespace Kromer.Models.Api.V1.Subscriptions;
+
+public class SubscriptionListResponse
+{
+ public int Count { get; set; }
+
+ public int Total { get; set; }
+
+ public IEnumerable Subscriptions { get; set; } = [];
+}
diff --git a/Kromer/Models/Dto/SubscriptionDto.cs b/Kromer/Models/Dto/SubscriptionDto.cs
new file mode 100644
index 0000000..c47d35f
--- /dev/null
+++ b/Kromer/Models/Dto/SubscriptionDto.cs
@@ -0,0 +1,33 @@
+using System.Text.Json.Serialization;
+using Kromer.Models.Entities;
+
+namespace Kromer.Models.Dto;
+
+public class SubscriptionDto
+{
+ public int Id { get; set; }
+
+ public string Description { get; set; } = null!;
+
+ public decimal Price { get; set; }
+
+ public int Period { get; set; }
+
+ public string Name { get; set; } = null!;
+
+ public int Subscribers { get; set; }
+
+ [JsonIgnore(Condition = JsonIgnoreCondition.WhenWritingNull)]
+ public int? MaxSubscribers { get; set; }
+
+ [JsonIgnore(Condition = JsonIgnoreCondition.WhenWritingNull)]
+ public IReadOnlyCollection? AllowedSubscribers { get; set; }
+
+ public SubscriptionStatus Status { get; set; }
+
+ [JsonIgnore(Condition = JsonIgnoreCondition.WhenWritingNull)]
+ public string? OwnerAddress { get; set; }
+
+ [JsonIgnore(Condition = JsonIgnoreCondition.WhenWritingNull)]
+ public IReadOnlyCollection? WalletSubscriptions { get; set; }
+}
diff --git a/Kromer/Models/Dto/WalletSubscriptionDto.cs b/Kromer/Models/Dto/WalletSubscriptionDto.cs
new file mode 100644
index 0000000..d44a706
--- /dev/null
+++ b/Kromer/Models/Dto/WalletSubscriptionDto.cs
@@ -0,0 +1,10 @@
+namespace Kromer.Models.Dto;
+
+public class WalletSubscriptionDto
+{
+ public string Address { get; set; } = null!;
+
+ public DateTime NextPayment { get; set; }
+
+ public bool Unsubscribable { get; set; }
+}
diff --git a/Kromer/Models/Entities/SubscriptionContractEntity.cs b/Kromer/Models/Entities/SubscriptionContractEntity.cs
new file mode 100644
index 0000000..7ebb4c6
--- /dev/null
+++ b/Kromer/Models/Entities/SubscriptionContractEntity.cs
@@ -0,0 +1,28 @@
+namespace Kromer.Models.Entities;
+
+public class SubscriptionContractEntity
+{
+ public int Id { get; set; }
+
+ public string Receiver { get; set; } = null!;
+
+ public string BaseName { get; set; } = null!;
+
+ public decimal Price { get; set; }
+
+ public int PeriodMinutes { get; set; }
+
+ public string Description { get; set; } = null!;
+
+ public int? MaxSubscribers { get; set; }
+
+ public string[]? AllowedSubscriberAddresses { get; set; }
+
+ public SubscriptionStatus Status { get; set; } = SubscriptionStatus.Active;
+
+ public DateTime CreatedAt { get; set; }
+
+ public DateTime? CancelledAt { get; set; }
+
+ public ICollection WalletSubscriptions { get; set; } = [];
+}
diff --git a/Kromer/Models/Entities/SubscriptionStatus.cs b/Kromer/Models/Entities/SubscriptionStatus.cs
new file mode 100644
index 0000000..d642acd
--- /dev/null
+++ b/Kromer/Models/Entities/SubscriptionStatus.cs
@@ -0,0 +1,8 @@
+namespace Kromer.Models.Entities;
+
+public enum SubscriptionStatus
+{
+ Active,
+ Closed,
+ Cancelled,
+}
diff --git a/Kromer/Models/Entities/WalletSubscriptionEntity.cs b/Kromer/Models/Entities/WalletSubscriptionEntity.cs
new file mode 100644
index 0000000..ce5059b
--- /dev/null
+++ b/Kromer/Models/Entities/WalletSubscriptionEntity.cs
@@ -0,0 +1,24 @@
+namespace Kromer.Models.Entities;
+
+public class WalletSubscriptionEntity
+{
+ public int Id { get; set; }
+
+ public int ContractId { get; set; }
+
+ public SubscriptionContractEntity Contract { get; set; } = null!;
+
+ public string WalletAddress { get; set; } = null!;
+
+ public DateTime NextPayment { get; set; }
+
+ public SubscriptionStatus Status { get; set; } = SubscriptionStatus.Active;
+
+ public string? CancellationReason { get; set; }
+
+ public bool CanUnsubscribe { get; set; } = true;
+
+ public DateTime CreatedAt { get; set; }
+
+ public DateTime? CancelledAt { get; set; }
+}
diff --git a/Kromer/Models/Exceptions/ErrorCode.cs b/Kromer/Models/Exceptions/ErrorCode.cs
index f60e48f..84a98ad 100644
--- a/Kromer/Models/Exceptions/ErrorCode.cs
+++ b/Kromer/Models/Exceptions/ErrorCode.cs
@@ -46,6 +46,26 @@ public enum ErrorCode
[StatusCode(HttpStatusCode.Forbidden)]
SameWalletTransfer,
+ [Description("The subscription contract is closed")]
+ [StatusCode(HttpStatusCode.Forbidden)]
+ SubscriptionClosed,
+
+ [Description("The subscription contract is cancelled")]
+ [StatusCode(HttpStatusCode.Gone)]
+ SubscriptionCancelled,
+
+ [Description("The subscription contract is full")]
+ [StatusCode(HttpStatusCode.Conflict)]
+ SubscriptionFull,
+
+ [Description("The wallet is not allowed to subscribe to this contract")]
+ [StatusCode(HttpStatusCode.Forbidden)]
+ SubscriberNotAllowed,
+
+ [Description("The subscription cannot be unsubscribed")]
+ [StatusCode(HttpStatusCode.Forbidden)]
+ SubscriptionCannotUnsubscribe,
+
[Description("Resource not found")]
[StatusCode(HttpStatusCode.NotFound)]
ResourceNotFound,
@@ -61,4 +81,4 @@ public enum ErrorCode
[Description("Invalid request type")]
[StatusCode(HttpStatusCode.BadRequest)]
InvalidRequestType,
-}
\ No newline at end of file
+}
diff --git a/Kromer/Models/WebSocket/Events/KromerSubscriptionEvent.cs b/Kromer/Models/WebSocket/Events/KromerSubscriptionEvent.cs
new file mode 100644
index 0000000..ca01851
--- /dev/null
+++ b/Kromer/Models/WebSocket/Events/KromerSubscriptionEvent.cs
@@ -0,0 +1,24 @@
+namespace Kromer.Models.WebSocket.Events;
+
+public class KromerSubscriptionEvent : IKristEvent
+{
+ public string Type => "event";
+
+ public string Event => "subscription";
+
+ public string Action { get; set; } = null!;
+
+ public int ContractId { get; set; }
+
+ public int? SubscriptionId { get; set; }
+
+ public string? OwnerAddress { get; set; }
+
+ public string? SubscriberAddress { get; set; }
+
+ public string Status { get; set; } = null!;
+
+ public string? Reason { get; set; }
+
+ public DateTime? NextPayment { get; set; }
+}
diff --git a/Kromer/Models/WebSocket/SubscriptionLevels.cs b/Kromer/Models/WebSocket/SubscriptionLevels.cs
index 7eb9d34..07fe58c 100644
--- a/Kromer/Models/WebSocket/SubscriptionLevels.cs
+++ b/Kromer/Models/WebSocket/SubscriptionLevels.cs
@@ -14,4 +14,6 @@ public enum SubscriptionLevel : byte
[Obsolete("Unused")]
Motd = 64,
-}
\ No newline at end of file
+
+ OwnSubscriptions = 128,
+}
diff --git a/Kromer/Program.cs b/Kromer/Program.cs
index 6775932..8342d70 100644
--- a/Kromer/Program.cs
+++ b/Kromer/Program.cs
@@ -23,7 +23,11 @@
builder.Services.AddDbContext(options =>
options.UseNpgsql(builder.Configuration.GetConnectionString("Default"),
- o => o.MapEnum("transaction_type", "public")));
+ o =>
+ {
+ o.MapEnum("transaction_type", "public");
+ o.MapEnum("subscription_status", "public");
+ }));
builder.Services.AddScoped();
builder.Services.AddScoped();
@@ -31,6 +35,7 @@
builder.Services.AddScoped();
builder.Services.AddScoped();
builder.Services.AddScoped();
+builder.Services.AddScoped();
builder.Services.AddScoped();
builder.Services.AddScoped();
@@ -43,6 +48,7 @@
builder.Services.AddHostedService();
builder.Services.AddHostedService();
+builder.Services.AddHostedService();
// Support for reverse proxies, like NGINX
builder.Services.Configure(options => { options.ForwardedHeaders = ForwardedHeaders.All; });
@@ -109,6 +115,32 @@
var app = builder.Build();
+// Please don't murder me for this, I think this is how one is supposed to auto apply migrations..
+using (var scope = app.Services.CreateScope())
+{
+ var db = scope.ServiceProvider.GetRequiredService();
+ var hasExistingKromerSchema = db.Database
+ .SqlQueryRaw("SELECT EXISTS (SELECT 1 FROM information_schema.tables WHERE table_schema = 'public' AND table_name = 'wallets') AS \"Value\"")
+ .Single();
+
+ if (hasExistingKromerSchema)
+ {
+ var pendingMigrations = db.Database.GetPendingMigrations().ToArray();
+
+ if (pendingMigrations.Length > 0)
+ {
+ app.Logger.LogInformation("Applying {MigrationCount} pending database migration(s): {Migrations}",
+ pendingMigrations.Length, string.Join(", ", pendingMigrations));
+
+ db.Database.Migrate();
+ }
+ }
+ else
+ {
+ app.Logger.LogWarning("Skipping database migrations because the existing Kromer schema was not found.");
+ }
+}
+
app.UseForwardedHeaders();
// Configure the HTTP request pipeline.
@@ -210,4 +242,4 @@ await context.Response.WriteAsJsonAsync(new KristResult
});
});
-app.Run();
\ No newline at end of file
+app.Run();
diff --git a/Kromer/Repositories/SubscriptionRepository.cs b/Kromer/Repositories/SubscriptionRepository.cs
new file mode 100644
index 0000000..e498d3d
--- /dev/null
+++ b/Kromer/Repositories/SubscriptionRepository.cs
@@ -0,0 +1,825 @@
+using System.Threading.Channels;
+using Kromer.Data;
+using Kromer.Models.Api.V1.Subscriptions;
+using Kromer.Models.Dto;
+using Kromer.Models.Entities;
+using Kromer.Models.Exceptions;
+using Kromer.Models.WebSocket.Events;
+using Kromer.Services;
+using Kromer.Utils;
+using Microsoft.EntityFrameworkCore;
+
+namespace Kromer.Repositories;
+
+public class SubscriptionRepository(
+ KromerContext context,
+ WalletRepository walletRepository,
+ TransactionService transactionService,
+ Channel eventChannel,
+ ILogger logger)
+{
+ private const string ReasonContractCancelled = "contract_cancelled";
+ private const string ReasonInsufficientFunds = "insufficient_funds";
+ private const string ReasonUnsubscribed = "unsubscribed";
+ private const string ReasonSubscriberMissing = "subscriber_missing";
+ private const string ReasonOwnerMissing = "owner_missing";
+ private const string ReasonSameWallet = "same_wallet";
+ private const int MaxAllowedSubscriberAddresses = 1000;
+ private const int MaxQueryValues = 1000;
+
+ public async Task CreateContractAsync(CreateSubscriptionRequest? request)
+ {
+ AssertCreateRequest(request);
+
+ var receiver = await NormalizeReceiverAsync(request!.Name, requireExisting: true);
+ var owner = await AuthenticateCurrentOwnerAsync(request.PrivateKey, receiver.BaseName);
+
+ var contract = new SubscriptionContractEntity
+ {
+ Receiver = receiver.Receiver,
+ BaseName = receiver.BaseName,
+ Price = decimal.Round(request.Price, 5, MidpointRounding.ToEven),
+ PeriodMinutes = request.Period,
+ Description = request.Description.Trim(),
+ MaxSubscribers = request.MaxSubscribers,
+ AllowedSubscriberAddresses = NormalizeAllowedSubscribers(request.AllowedSubscribers),
+ Status = SubscriptionStatus.Active,
+ CreatedAt = DateTime.UtcNow,
+ };
+
+ await context.SubscriptionContracts.AddAsync(contract);
+ await context.SaveChangesAsync();
+
+ logger.LogInformation("Created subscription contract {ContractId} for {Receiver} owned by {Owner}",
+ contract.Id, contract.Receiver, owner.Address);
+
+ return new CreateSubscriptionResponse
+ {
+ Id = contract.Id,
+ };
+ }
+
+ public async Task CloseContractAsync(int id, string? privateKey)
+ {
+ AssertPrivateKey(privateKey);
+
+ var contract = await context.SubscriptionContracts.FirstOrDefaultAsync(q => q.Id == id);
+ if (contract is null)
+ {
+ throw new KromerException(ErrorCode.ResourceNotFound);
+ }
+
+ var owner = await AuthenticateCurrentOwnerAsync(privateKey, contract.BaseName);
+
+ if (contract.Status == SubscriptionStatus.Active)
+ {
+ contract.Status = SubscriptionStatus.Closed;
+ await context.SaveChangesAsync();
+
+ await EmitSubscriptionEventAsync("contract_closed", contract, null, owner.Address,
+ SubscriptionStatus.Closed);
+ }
+
+ return await BuildDtoAsync(contract, [owner.Address]);
+ }
+
+ public async Task CancelContractAsync(int id, string? privateKey)
+ {
+ AssertPrivateKey(privateKey);
+
+ var now = DateTime.UtcNow;
+ var contract = await context.SubscriptionContracts
+ .Include(q => q.WalletSubscriptions.Where(s =>
+ s.Status == SubscriptionStatus.Active ||
+ (s.CancellationReason == ReasonUnsubscribed && s.NextPayment > now)))
+ .FirstOrDefaultAsync(q => q.Id == id);
+
+ if (contract is null)
+ {
+ throw new KromerException(ErrorCode.ResourceNotFound);
+ }
+
+ var owner = await AuthenticateCurrentOwnerAsync(privateKey, contract.BaseName);
+
+ if (contract.Status != SubscriptionStatus.Cancelled)
+ {
+ contract.Status = SubscriptionStatus.Cancelled;
+ contract.CancelledAt = now;
+
+ foreach (var subscription in contract.WalletSubscriptions)
+ {
+ CancelWalletSubscription(subscription, ReasonContractCancelled, now);
+ }
+
+ await context.SaveChangesAsync();
+
+ if (contract.WalletSubscriptions.Count == 0)
+ {
+ await EmitSubscriptionEventAsync("contract_cancelled", contract, null, owner.Address,
+ SubscriptionStatus.Cancelled, ReasonContractCancelled);
+ }
+ else
+ {
+ foreach (var subscription in contract.WalletSubscriptions)
+ {
+ await EmitSubscriptionEventAsync("contract_cancelled", contract, subscription, owner.Address,
+ SubscriptionStatus.Cancelled, ReasonContractCancelled);
+ }
+ }
+ }
+
+ return await BuildDtoAsync(contract, [owner.Address]);
+ }
+
+ public async Task GetContractAsync(int id, IEnumerable? addresses)
+ {
+ var contextAddresses = NormalizeAddressFilters(addresses);
+
+ var contract = await context.SubscriptionContracts.FirstOrDefaultAsync(q => q.Id == id);
+ if (contract is null)
+ {
+ throw new KromerException(ErrorCode.ResourceNotFound);
+ }
+
+ return await BuildDtoAsync(contract, contextAddresses);
+ }
+
+ public async Task ListContractsAsync(
+ IEnumerable? addresses,
+ IEnumerable? names,
+ bool excludeOwned,
+ bool onlyOwned,
+ bool onlyUnsubscribable,
+ int limit,
+ int offset)
+ {
+ limit = Math.Clamp(limit, 1, 1000);
+ offset = Math.Max(offset, 0);
+
+ var query = context.SubscriptionContracts
+ .Where(q => q.Status != SubscriptionStatus.Cancelled);
+
+ var contextAddresses = NormalizeAddressFilters(addresses);
+ var receiverFilters = await NormalizeReceiverFiltersAsync(names, requireExisting: false);
+ if (receiverFilters.Count > 0)
+ {
+ var receivers = receiverFilters
+ .Where(q => q.IsMetaname)
+ .Select(q => q.Receiver)
+ .ToArray();
+ var baseNames = receiverFilters
+ .Where(q => !q.IsMetaname)
+ .Select(q => q.BaseName)
+ .ToArray();
+
+ query = query.Where(q => receivers.Contains(q.Receiver) || baseNames.Contains(q.BaseName));
+ }
+
+ if (contextAddresses.Count > 0)
+ {
+ var now = DateTime.UtcNow;
+ var ownedBaseNames = context.Names
+ .Where(q => contextAddresses.Contains(q.Owner))
+ .Select(q => q.Name);
+
+ var subscribedContracts = context.WalletSubscriptions
+ .Where(q => contextAddresses.Contains(q.WalletAddress) &&
+ (q.Status == SubscriptionStatus.Active ||
+ (q.CancellationReason == ReasonUnsubscribed && q.NextPayment > now)));
+
+ if (onlyUnsubscribable)
+ {
+ subscribedContracts = subscribedContracts.Where(q =>
+ q.Status == SubscriptionStatus.Active && q.CanUnsubscribe);
+ }
+
+ var subscribedIds = subscribedContracts.Select(q => q.ContractId);
+
+ if (onlyOwned)
+ {
+ query = query.Where(q => ownedBaseNames.Contains(q.BaseName));
+ }
+ else if (excludeOwned)
+ {
+ query = query.Where(q => subscribedIds.Contains(q.Id) && !ownedBaseNames.Contains(q.BaseName));
+ }
+ else if (receiverFilters.Count == 0)
+ {
+ query = query.Where(q => subscribedIds.Contains(q.Id) || ownedBaseNames.Contains(q.BaseName));
+ }
+ }
+
+ if (receiverFilters.Count == 0 && contextAddresses.Count == 0)
+ {
+ throw new KromerException(ErrorCode.InvalidParameter);
+ }
+
+ var total = await query.CountAsync();
+ var contracts = await query
+ .OrderBy(q => q.Id)
+ .Skip(offset)
+ .Take(limit)
+ .ToListAsync();
+
+ var dtos = new List();
+ foreach (var contract in contracts)
+ {
+ dtos.Add(await BuildDtoAsync(contract, contextAddresses));
+ }
+
+ return new SubscriptionListResponse
+ {
+ Count = dtos.Count,
+ Total = total,
+ Subscriptions = dtos,
+ };
+ }
+
+ public async Task SubscribeAsync(int contractId, string? privateKey)
+ {
+ AssertPrivateKey(privateKey);
+
+ var subscriber = await walletRepository.GetWalletFromKeyAsync(privateKey!);
+ if (subscriber is null)
+ {
+ throw new KromerException(ErrorCode.AuthenticationFailed);
+ }
+
+ var now = DateTime.UtcNow;
+ SubscriptionContractEntity contract;
+ WalletSubscriptionEntity subscription;
+ WalletEntity owner;
+
+ await using (var transaction = await context.Database.BeginTransactionAsync())
+ {
+ var lockedContract = await context.SubscriptionContracts
+ .FromSqlInterpolated(
+ $"SELECT * FROM subscription_contracts WHERE id = {contractId} FOR UPDATE")
+ .FirstOrDefaultAsync();
+ if (lockedContract is null)
+ {
+ throw new KromerException(ErrorCode.ResourceNotFound);
+ }
+
+ contract = lockedContract;
+ if (contract.Status == SubscriptionStatus.Cancelled)
+ {
+ throw new KromerException(ErrorCode.SubscriptionCancelled);
+ }
+
+ var existing = await context.WalletSubscriptions.FirstOrDefaultAsync(q =>
+ q.ContractId == contract.Id &&
+ q.WalletAddress == subscriber.Address &&
+ q.Status == SubscriptionStatus.Active);
+ if (existing is not null)
+ {
+ await transaction.CommitAsync();
+ return new SubscribeResponse
+ {
+ NextPayment = existing.NextPayment,
+ };
+ }
+
+ owner = await ResolveCurrentOwnerWalletAsync(contract.BaseName);
+ if (owner.Address == subscriber.Address)
+ {
+ throw new KromerException(ErrorCode.SameWalletTransfer);
+ }
+
+ var cancelledWithRemainingTime = await context.WalletSubscriptions.FirstOrDefaultAsync(q =>
+ q.ContractId == contract.Id &&
+ q.WalletAddress == subscriber.Address &&
+ q.CancellationReason == ReasonUnsubscribed &&
+ q.NextPayment > now);
+ if (cancelledWithRemainingTime is not null)
+ {
+ cancelledWithRemainingTime.Status = SubscriptionStatus.Active;
+ cancelledWithRemainingTime.CancellationReason = null;
+ cancelledWithRemainingTime.CancelledAt = null;
+ context.WalletSubscriptions.Update(cancelledWithRemainingTime);
+ await context.SaveChangesAsync();
+ await transaction.CommitAsync();
+
+ await EmitSubscriptionEventAsync("subscribe", contract, cancelledWithRemainingTime, owner.Address,
+ SubscriptionStatus.Active);
+
+ return new SubscribeResponse
+ {
+ NextPayment = cancelledWithRemainingTime.NextPayment,
+ };
+ }
+
+ if (contract.Status == SubscriptionStatus.Closed)
+ {
+ throw new KromerException(ErrorCode.SubscriptionClosed);
+ }
+
+ AssertSubscriberAllowed(contract, subscriber.Address);
+ await AssertSubscriberCapacityAsync(contract, now);
+
+ subscription = new WalletSubscriptionEntity
+ {
+ ContractId = contract.Id,
+ WalletAddress = subscriber.Address,
+ NextPayment = now.AddMinutes(contract.PeriodMinutes),
+ Status = SubscriptionStatus.Active,
+ CanUnsubscribe = true,
+ CreatedAt = now,
+ };
+
+ await context.WalletSubscriptions.AddAsync(subscription);
+ await context.SaveChangesAsync();
+ await transaction.CommitAsync();
+ }
+
+ try
+ {
+ await ChargeSubscriptionAsync(contract, subscription, subscriber, owner, subscription.NextPayment, now);
+ }
+ catch (KristException ex) when (ex.Code == ErrorCode.InsufficientFunds)
+ {
+ CancelWalletSubscription(subscription, ReasonInsufficientFunds, DateTime.UtcNow);
+ await context.SaveChangesAsync();
+ await EmitSubscriptionEventAsync("payment_failed", contract, subscription, owner.Address,
+ SubscriptionStatus.Cancelled, ReasonInsufficientFunds);
+ throw;
+ }
+
+ await EmitSubscriptionEventAsync("subscribe", contract, subscription, owner.Address, SubscriptionStatus.Active);
+
+ return new SubscribeResponse
+ {
+ NextPayment = subscription.NextPayment,
+ };
+ }
+
+ public async Task UnsubscribeAsync(int contractId, string? privateKey)
+ {
+ AssertPrivateKey(privateKey);
+
+ var subscriber = await walletRepository.GetWalletFromKeyAsync(privateKey!);
+ if (subscriber is null)
+ {
+ throw new KromerException(ErrorCode.AuthenticationFailed);
+ }
+
+ var contract = await context.SubscriptionContracts.FirstOrDefaultAsync(q => q.Id == contractId);
+ if (contract is null)
+ {
+ throw new KromerException(ErrorCode.ResourceNotFound);
+ }
+
+ var subscription = await context.WalletSubscriptions.FirstOrDefaultAsync(q =>
+ q.ContractId == contract.Id &&
+ q.WalletAddress == subscriber.Address &&
+ q.Status == SubscriptionStatus.Active);
+ if (subscription is null)
+ {
+ return;
+ }
+
+ if (!subscription.CanUnsubscribe)
+ {
+ throw new KromerException(ErrorCode.SubscriptionCannotUnsubscribe);
+ }
+
+ CancelWalletSubscription(subscription, ReasonUnsubscribed, DateTime.UtcNow);
+ await context.SaveChangesAsync();
+
+ var ownerAddress = await ResolveCurrentOwnerAddressAsync(contract.BaseName, throwIfMissing: false);
+ await EmitSubscriptionEventAsync("unsubscribe", contract, subscription, ownerAddress,
+ SubscriptionStatus.Cancelled, ReasonUnsubscribed);
+ }
+
+ public async Task BillDueSubscriptionsAsync(CancellationToken cancellationToken)
+ {
+ var now = DateTime.UtcNow;
+ var subscriptions = await context.WalletSubscriptions
+ .Include(q => q.Contract)
+ .Where(q =>
+ q.Status == SubscriptionStatus.Active &&
+ q.Contract.Status != SubscriptionStatus.Cancelled &&
+ q.NextPayment <= now)
+ .OrderBy(q => q.NextPayment)
+ .Take(100)
+ .ToListAsync(cancellationToken);
+
+ foreach (var subscription in subscriptions)
+ {
+ try
+ {
+ await BillSubscriptionAsync(subscription, now, cancellationToken);
+ }
+ catch (Exception ex)
+ {
+ logger.LogError(ex, "Failed to process subscription {SubscriptionId} for contract {ContractId}",
+ subscription.Id, subscription.ContractId);
+ }
+ }
+
+ return subscriptions.Count;
+ }
+
+ private async Task BillSubscriptionAsync(WalletSubscriptionEntity subscription, DateTime now,
+ CancellationToken cancellationToken)
+ {
+ var contract = subscription.Contract;
+ var ownerAddress = await ResolveCurrentOwnerAddressAsync(contract.BaseName, throwIfMissing: false);
+ if (ownerAddress is null)
+ {
+ await CancelDueSubscriptionAsync(contract, subscription, null, ReasonOwnerMissing, cancellationToken);
+ return;
+ }
+
+ var owner = await walletRepository.GetWalletFromAddress(ownerAddress);
+ if (owner is null)
+ {
+ await CancelDueSubscriptionAsync(contract, subscription, ownerAddress, ReasonOwnerMissing,
+ cancellationToken);
+ return;
+ }
+
+ var subscriber = await walletRepository.GetWalletFromAddress(subscription.WalletAddress);
+ if (subscriber is null)
+ {
+ await CancelDueSubscriptionAsync(contract, subscription, ownerAddress, ReasonSubscriberMissing,
+ cancellationToken);
+ return;
+ }
+
+ if (subscriber.Address == owner.Address)
+ {
+ await CancelDueSubscriptionAsync(contract, subscription, ownerAddress, ReasonSameWallet,
+ cancellationToken);
+ return;
+ }
+
+ while (subscription.Status == SubscriptionStatus.Active && subscription.NextPayment <= now)
+ {
+ var dueAt = subscription.NextPayment;
+ var nextPayment = dueAt.AddMinutes(contract.PeriodMinutes);
+
+ try
+ {
+ await ChargeSubscriptionAsync(contract, subscription, subscriber, owner, nextPayment, dueAt);
+ }
+ catch (KristException ex) when (ex.Code == ErrorCode.InsufficientFunds)
+ {
+ await CancelDueSubscriptionAsync(contract, subscription, ownerAddress, ReasonInsufficientFunds,
+ cancellationToken);
+ }
+ }
+ }
+
+ private async Task CancelDueSubscriptionAsync(SubscriptionContractEntity contract,
+ WalletSubscriptionEntity subscription, string? ownerAddress, string reason, CancellationToken cancellationToken)
+ {
+ CancelWalletSubscription(subscription, reason, DateTime.UtcNow);
+ await context.SaveChangesAsync(cancellationToken);
+ await EmitSubscriptionEventAsync("payment_failed", contract, subscription, ownerAddress,
+ SubscriptionStatus.Cancelled, reason);
+ }
+
+ private async Task ChargeSubscriptionAsync(SubscriptionContractEntity contract, WalletSubscriptionEntity subscription,
+ WalletEntity subscriber, WalletEntity owner, DateTime nextPayment, DateTime transactionDate)
+ {
+ subscription.NextPayment = nextPayment;
+ context.WalletSubscriptions.Update(subscription);
+
+ var transaction = transactionService.InitiateTransaction(subscriber, owner, contract.Price);
+ transaction.Date = transactionDate;
+ transaction.Metadata = $"subscription={contract.Id};wallet_subscription={subscription.Id}";
+
+ var receiver = Validation.ParseMetaName(contract.Receiver);
+ if (receiver.Valid)
+ {
+ transaction.SentName = receiver.Name;
+ transaction.SentMetaname = string.IsNullOrWhiteSpace(receiver.Meta) ? null : receiver.Meta;
+ }
+
+ await transactionService.CommitTransactionAsync(subscriber, owner, transaction);
+
+ await eventChannel.Writer.WriteAsync(new KristTransactionEvent
+ {
+ Transaction = TransactionDto.FromEntity(transaction),
+ });
+
+ await EmitSubscriptionEventAsync("payment_success", contract, subscription, owner.Address,
+ SubscriptionStatus.Active);
+ }
+
+ private async Task BuildDtoAsync(SubscriptionContractEntity contract,
+ IReadOnlyCollection contextAddresses)
+ {
+ var now = DateTime.UtcNow;
+ var subscribers = contract.Status == SubscriptionStatus.Cancelled
+ ? 0
+ : await CountCurrentSubscribersAsync(contract.Id, now);
+
+ var dto = new SubscriptionDto
+ {
+ Id = contract.Id,
+ Description = contract.Description,
+ Price = contract.Price,
+ Period = contract.PeriodMinutes,
+ Name = contract.Receiver,
+ Subscribers = subscribers,
+ MaxSubscribers = contract.MaxSubscribers,
+ AllowedSubscribers = contract.AllowedSubscriberAddresses,
+ Status = contract.Status,
+ };
+
+ dto.OwnerAddress = await ResolveCurrentOwnerAddressAsync(contract.BaseName, throwIfMissing: false);
+
+ if (contextAddresses.Count == 0)
+ {
+ return dto;
+ }
+
+ var contextAddressArray = contextAddresses.ToArray();
+ var subscriptions = await context.WalletSubscriptions
+ .Where(q =>
+ q.ContractId == contract.Id &&
+ contextAddressArray.Contains(q.WalletAddress) &&
+ contract.Status != SubscriptionStatus.Cancelled &&
+ (q.Status == SubscriptionStatus.Active ||
+ (q.CancellationReason == ReasonUnsubscribed && q.NextPayment > now)))
+ .OrderBy(q => q.WalletAddress)
+ .Select(q => new WalletSubscriptionDto
+ {
+ Address = q.WalletAddress,
+ NextPayment = q.NextPayment,
+ Unsubscribable = q.Status == SubscriptionStatus.Active && q.CanUnsubscribe,
+ })
+ .ToListAsync();
+
+ dto.WalletSubscriptions = subscriptions.Count > 0 ? subscriptions : null;
+
+ return dto;
+ }
+
+ private async Task AuthenticateCurrentOwnerAsync(string? privateKey, string baseName)
+ {
+ AssertPrivateKey(privateKey);
+
+ var wallet = await walletRepository.GetWalletFromKeyAsync(privateKey!);
+ if (wallet is null)
+ {
+ throw new KromerException(ErrorCode.AuthenticationFailed);
+ }
+
+ var ownerAddress = await ResolveCurrentOwnerAddressAsync(baseName, throwIfMissing: true);
+ if (ownerAddress != wallet.Address)
+ {
+ throw new KromerException(ErrorCode.NotNameOwner);
+ }
+
+ return wallet;
+ }
+
+ private async Task ResolveCurrentOwnerWalletAsync(string baseName)
+ {
+ var ownerAddress = await ResolveCurrentOwnerAddressAsync(baseName, throwIfMissing: true);
+ var wallet = await walletRepository.GetWalletFromAddress(ownerAddress!);
+ if (wallet is null)
+ {
+ throw new KromerException(ErrorCode.AddressNotFound);
+ }
+
+ return wallet;
+ }
+
+ private async Task ResolveCurrentOwnerAddressAsync(string baseName, bool throwIfMissing)
+ {
+ var name = await context.Names.FirstOrDefaultAsync(q => q.Name == baseName);
+ if (name is not null)
+ {
+ return name.Owner;
+ }
+
+ if (throwIfMissing)
+ {
+ throw new KromerException(ErrorCode.NameNotFound);
+ }
+
+ return null;
+ }
+
+ private async Task NormalizeReceiverAsync(string? value, bool requireExisting)
+ {
+ if (string.IsNullOrWhiteSpace(value))
+ {
+ throw new KromerException(ErrorCode.InvalidParameter);
+ }
+
+ value = value.Trim().ToLowerInvariant();
+ if (Validation.IsValidAddress(value))
+ {
+ throw new KromerException(ErrorCode.InvalidParameter);
+ }
+
+ string receiver;
+ string baseName;
+ var isMetaname = false;
+
+ if (value.EndsWith(".kro"))
+ {
+ var parsed = Validation.ParseMetaName(value);
+ if (!parsed.Valid || !Validation.IsNameValid(parsed.Name))
+ {
+ throw new KromerException(ErrorCode.InvalidParameter);
+ }
+
+ baseName = Validation.SanitizeName(parsed.Name);
+ var meta = parsed.Meta?.Trim().ToLowerInvariant();
+ isMetaname = !string.IsNullOrWhiteSpace(meta);
+ receiver = isMetaname ? $"{meta}@{baseName}.kro" : $"{baseName}.kro";
+ }
+ else
+ {
+ if (value.Contains('@') || !Validation.IsNameValid(value))
+ {
+ throw new KromerException(ErrorCode.InvalidParameter);
+ }
+
+ baseName = Validation.SanitizeName(value);
+ receiver = $"{baseName}.kro";
+ }
+
+ if (requireExisting && !await context.Names.AnyAsync(q => q.Name == baseName))
+ {
+ throw new KromerException(ErrorCode.NameNotFound);
+ }
+
+ return new NormalizedReceiver(receiver, baseName, isMetaname);
+ }
+
+ private static List NormalizeAddressFilters(IEnumerable? values)
+ {
+ if (values is null)
+ {
+ return [];
+ }
+
+ var normalized = values
+ .Select(q => q.Trim().ToLowerInvariant())
+ .ToList();
+
+ if (normalized.Count > MaxQueryValues ||
+ normalized.Any(q => string.IsNullOrWhiteSpace(q) || !Validation.IsValidAddress(q)))
+ {
+ throw new KromerException(ErrorCode.InvalidParameter);
+ }
+
+ return normalized.Distinct(StringComparer.Ordinal).ToList();
+ }
+
+ private async Task> NormalizeReceiverFiltersAsync(IEnumerable? values,
+ bool requireExisting)
+ {
+ var receivers = new List();
+ var seen = new HashSet(StringComparer.Ordinal);
+
+ if (values is null)
+ {
+ return receivers;
+ }
+
+ var queryValues = values.ToList();
+ if (queryValues.Count > MaxQueryValues)
+ {
+ throw new KromerException(ErrorCode.InvalidParameter);
+ }
+
+ foreach (var value in queryValues)
+ {
+ var receiver = await NormalizeReceiverAsync(value.Trim(), requireExisting);
+ var key = receiver.IsMetaname ? $"receiver:{receiver.Receiver}" : $"base:{receiver.BaseName}";
+
+ if (seen.Add(key))
+ {
+ receivers.Add(receiver);
+ }
+ }
+
+ return receivers;
+ }
+
+ private static void AssertCreateRequest(CreateSubscriptionRequest? request)
+ {
+ if (request is null ||
+ string.IsNullOrWhiteSpace(request.PrivateKey) ||
+ string.IsNullOrWhiteSpace(request.Name) ||
+ string.IsNullOrWhiteSpace(request.Description) ||
+ request.Description.Trim().Length > 255 ||
+ request.MaxSubscribers is <= 0 ||
+ request.Period < 1 ||
+ decimal.Round(request.Price, 5, MidpointRounding.ToEven) <= 0)
+ {
+ throw new KromerException(ErrorCode.InvalidParameter);
+ }
+ }
+
+ private static void AssertPrivateKey(string? privateKey)
+ {
+ if (string.IsNullOrWhiteSpace(privateKey))
+ {
+ throw new KromerException(ErrorCode.InvalidParameter);
+ }
+ }
+
+ private static void CancelWalletSubscription(WalletSubscriptionEntity subscription, string reason, DateTime now)
+ {
+ subscription.Status = SubscriptionStatus.Cancelled;
+ subscription.CancellationReason = reason;
+ subscription.CancelledAt = now;
+ }
+
+ private async Task AssertSubscriberCapacityAsync(SubscriptionContractEntity contract, DateTime now)
+ {
+ if (contract.MaxSubscribers is null)
+ {
+ return;
+ }
+
+ var subscribers = await CountCurrentSubscribersAsync(contract.Id, now);
+ if (subscribers >= contract.MaxSubscribers.Value)
+ {
+ throw new KromerException(ErrorCode.SubscriptionFull);
+ }
+ }
+
+ private async Task CountCurrentSubscribersAsync(int contractId, DateTime now)
+ {
+ return await context.WalletSubscriptions.CountAsync(q =>
+ q.ContractId == contractId &&
+ (q.Status == SubscriptionStatus.Active ||
+ (q.CancellationReason == ReasonUnsubscribed && q.NextPayment > now)));
+ }
+
+ private static void AssertSubscriberAllowed(SubscriptionContractEntity contract, string subscriberAddress)
+ {
+ if (contract.AllowedSubscriberAddresses is not { Length: > 0 })
+ {
+ return;
+ }
+
+ if (!contract.AllowedSubscriberAddresses.Contains(subscriberAddress, StringComparer.Ordinal))
+ {
+ throw new KromerException(ErrorCode.SubscriberNotAllowed);
+ }
+ }
+
+ private static string[]? NormalizeAllowedSubscribers(IEnumerable? addresses)
+ {
+ if (addresses is null)
+ {
+ return null;
+ }
+
+ var normalized = new List();
+ foreach (var address in addresses)
+ {
+ if (string.IsNullOrWhiteSpace(address))
+ {
+ throw new KromerException(ErrorCode.InvalidParameter);
+ }
+
+ normalized.Add(address.Trim().ToLowerInvariant());
+ }
+
+ if (normalized.Count is 0 or > MaxAllowedSubscriberAddresses ||
+ normalized.Any(q => !Validation.IsValidAddress(q)))
+ {
+ throw new KromerException(ErrorCode.InvalidParameter);
+ }
+
+ return normalized
+ .Distinct(StringComparer.Ordinal)
+ .OrderBy(q => q, StringComparer.Ordinal)
+ .ToArray();
+ }
+
+ private async Task EmitSubscriptionEventAsync(string action, SubscriptionContractEntity contract,
+ WalletSubscriptionEntity? subscription, string? ownerAddress, SubscriptionStatus status, string? reason = null)
+ {
+ await eventChannel.Writer.WriteAsync(new KromerSubscriptionEvent
+ {
+ Action = action,
+ ContractId = contract.Id,
+ SubscriptionId = subscription?.Id,
+ OwnerAddress = ownerAddress,
+ SubscriberAddress = subscription?.WalletAddress,
+ Status = SnakeCaseNamingPolicy.Convert(status.ToString()),
+ Reason = reason,
+ NextPayment = subscription is not null &&
+ (subscription.Status == SubscriptionStatus.Active ||
+ (subscription.CancellationReason == ReasonUnsubscribed &&
+ subscription.NextPayment > DateTime.UtcNow))
+ ? subscription.NextPayment
+ : null,
+ });
+ }
+
+ private sealed record NormalizedReceiver(string Receiver, string BaseName, bool IsMetaname);
+}
diff --git a/Kromer/Services/SubscriptionBillingService.cs b/Kromer/Services/SubscriptionBillingService.cs
new file mode 100644
index 0000000..dbafe1b
--- /dev/null
+++ b/Kromer/Services/SubscriptionBillingService.cs
@@ -0,0 +1,49 @@
+using Kromer.Repositories;
+
+namespace Kromer.Services;
+
+public class SubscriptionBillingService(
+ IServiceScopeFactory scopeFactory,
+ ILogger logger) : BackgroundService
+{
+ protected override async Task ExecuteAsync(CancellationToken stoppingToken)
+ {
+ while (!stoppingToken.IsCancellationRequested)
+ {
+ try
+ {
+ await using var scope = scopeFactory.CreateAsyncScope();
+ var repository = scope.ServiceProvider.GetRequiredService();
+ var processed = 0;
+ var totalProcessed = 0;
+
+ do
+ {
+ processed = await repository.BillDueSubscriptionsAsync(stoppingToken);
+ totalProcessed += processed;
+
+ if (processed == 100)
+ {
+ await Task.Yield();
+ }
+ } while (processed == 100 &&
+ !stoppingToken.IsCancellationRequested);
+
+ if (totalProcessed > 0)
+ {
+ logger.LogInformation("Processed {Count} due subscription payments", totalProcessed);
+ }
+ }
+ catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
+ {
+ break;
+ }
+ catch (Exception ex)
+ {
+ logger.LogError(ex, "Failed to process due subscription payments");
+ }
+
+ await Task.Delay(TimeSpan.FromMinutes(1), stoppingToken);
+ }
+ }
+}
diff --git a/Kromer/SessionManager/EventDispatcher.cs b/Kromer/SessionManager/EventDispatcher.cs
index cc6ce97..eed3086 100644
--- a/Kromer/SessionManager/EventDispatcher.cs
+++ b/Kromer/SessionManager/EventDispatcher.cs
@@ -38,6 +38,13 @@ await Parallel.ForEachAsync(sessions, stoppingToken, async (session, token) =>
level = isOwn ? SubscriptionLevel.OwnTransactions : SubscriptionLevel.Transactions;
break;
}
+ case KromerSubscriptionEvent subscriptionEvent:
+ {
+ var isOwn = subscriptionEvent.OwnerAddress == address ||
+ subscriptionEvent.SubscriberAddress == address;
+ level = isOwn ? SubscriptionLevel.OwnSubscriptions : 0;
+ break;
+ }
}
if (level != 0 && session.SubscriptionLevel.HasFlag(level))
@@ -55,4 +62,4 @@ await Parallel.ForEachAsync(sessions, stoppingToken, async (session, token) =>
}
}
}
-}
\ No newline at end of file
+}
diff --git a/Kromer/appsettings.json b/Kromer/appsettings.json
index 43aa530..b4c3b89 100644
--- a/Kromer/appsettings.json
+++ b/Kromer/appsettings.json
@@ -1,12 +1,12 @@
{
"NameCost": 500,
"InitialBalance": 100.0,
- "InternalKey": "",
+ "InternalKey": "MEOWMEOW",
"PublicUrl": "localhost:5241",
"PublicWsUrl": "localhost:5241",
"AlertWebhook": "",
"ConnectionStrings": {
- "Default": ""
+ "Default": "Host=test-postgres;Port=5432;Database=testdb;Username=test;Password=test"
},
"Logging": {
"LogLevel": {