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
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@
using Microsoft.Extensions.AI;
using Netclaw.Actors.Channels;
using Netclaw.Actors.Protocol;
using Netclaw.Actors.Reminders;
using Netclaw.Configuration;
using Xunit;
using static Netclaw.Actors.Sessions.SessionProtocol;
Expand Down Expand Up @@ -91,7 +92,7 @@ await Source.Single(cmd)
SenderId = new SenderId("user-1"),
Contents = [new TextContent("hello")],
ReceivedAt = DateTimeOffset.UtcNow,
ReminderId = reminderId,
ReminderId = reminderId is null ? null : new ReminderId(reminderId),
AckTarget = ackTarget,
Audience = TrustAudience.Public,
Boundary = TrustBoundary.Public,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -276,7 +276,7 @@ public async Task Reminder_delivery_reports_success_when_post_succeeds()
var pipeline = new RecordingSessionPipeline(_ =>
[
new TextOutput("reminder output") { SessionId = sid },
new TurnCompleted { SessionId = sid, TurnNumber = new Netclaw.Actors.Protocol.TurnNumber(1), SourceReminderId = reminderKey }
new TurnCompleted { SessionId = sid, TurnNumber = new Netclaw.Actors.Protocol.TurnNumber(1), SourceReminderId = new ReminderId(reminderKey) }
], reactive: true);

var observer = CreateTestProbe();
Expand All @@ -286,7 +286,7 @@ public async Task Reminder_delivery_reports_success_when_post_succeeds()

var result = await observer.ExpectMsgAsync<ReminderDeliveryResult>(
TimeSpan.FromSeconds(5), cancellationToken: ct);
Assert.Equal(reminderKey, result.ReminderDeliveryKey);
Assert.Equal(new ReminderId(reminderKey), result.ReminderDeliveryKey);
Assert.Equal(ExpectedChannelType, result.ChannelType);
Assert.True(result.Delivered);
}
Expand All @@ -305,7 +305,7 @@ public async Task Reminder_delivery_reports_failure_when_post_throws()
var pipeline = new RecordingSessionPipeline(_ =>
[
new TextOutput("reminder output") { SessionId = sid },
new TurnCompleted { SessionId = sid, TurnNumber = new Netclaw.Actors.Protocol.TurnNumber(1), SourceReminderId = reminderKey }
new TurnCompleted { SessionId = sid, TurnNumber = new Netclaw.Actors.Protocol.TurnNumber(1), SourceReminderId = new ReminderId(reminderKey) }
], reactive: true);

SetReplyClientThrows(new InvalidOperationException("channel API down"));
Expand All @@ -316,7 +316,7 @@ public async Task Reminder_delivery_reports_failure_when_post_throws()

var result = await observer.ExpectMsgAsync<ReminderDeliveryResult>(
TimeSpan.FromSeconds(5), cancellationToken: ct);
Assert.Equal(reminderKey, result.ReminderDeliveryKey);
Assert.Equal(new ReminderId(reminderKey), result.ReminderDeliveryKey);
Assert.Equal(ExpectedChannelType, result.ChannelType);
Assert.False(result.Delivered);

Expand All @@ -339,9 +339,9 @@ public async Task Concurrent_reminders_to_same_session_each_get_their_own_result
var pipeline = new RecordingSessionPipeline(_ =>
[
new TextOutput("reply A") { SessionId = sid },
new TurnCompleted { SessionId = sid, TurnNumber = new Netclaw.Actors.Protocol.TurnNumber(1), SourceReminderId = keyA },
new TurnCompleted { SessionId = sid, TurnNumber = new Netclaw.Actors.Protocol.TurnNumber(1), SourceReminderId = new ReminderId(keyA) },
new TextOutput("reply B") { SessionId = sid },
new TurnCompleted { SessionId = sid, TurnNumber = new Netclaw.Actors.Protocol.TurnNumber(2), SourceReminderId = keyB }
new TurnCompleted { SessionId = sid, TurnNumber = new Netclaw.Actors.Protocol.TurnNumber(2), SourceReminderId = new ReminderId(keyB) }
], reactive: true);

var observerA = CreateTestProbe();
Expand All @@ -355,11 +355,11 @@ public async Task Concurrent_reminders_to_same_session_each_get_their_own_result

var resultA = await observerA.ExpectMsgAsync<ReminderDeliveryResult>(
TimeSpan.FromSeconds(5), cancellationToken: ct);
Assert.Equal(keyA, resultA.ReminderDeliveryKey);
Assert.Equal(new ReminderId(keyA), resultA.ReminderDeliveryKey);

var resultB = await observerB.ExpectMsgAsync<ReminderDeliveryResult>(
TimeSpan.FromSeconds(5), cancellationToken: ct);
Assert.Equal(keyB, resultB.ReminderDeliveryKey);
Assert.Equal(new ReminderId(keyB), resultB.ReminderDeliveryKey);
}

// Regression for the misleading-fallback bug: when the real content post
Expand Down Expand Up @@ -412,7 +412,7 @@ private MessageSource CreateReminderSource(string reminderKey, IActorRef deliver
SourceKind = new SourceKind("reminder")
},
ReceivedAt = DateTimeOffset.UnixEpoch,
ReminderId = reminderKey,
ReminderId = new ReminderId(reminderKey),
DeliveryObserver = deliveryObserver
};

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@
using Microsoft.Extensions.Hosting;
using Netclaw.Actors.Channels;
using Netclaw.Actors.Protocol;
using Netclaw.Actors.Reminders;
using Netclaw.Actors.Tests.Channels.TestHelpers;
using Netclaw.Channels.Discord;
using Netclaw.Configuration;
Expand Down Expand Up @@ -543,6 +544,6 @@ private static DiscordGatewayMessage CreateMessage(
{
SourceKind = new Netclaw.Actors.Channels.SourceKind("reminder")
},
ReminderId = "rem-1"
ReminderId = new ReminderId("rem-1")
};
}
5 changes: 3 additions & 2 deletions src/Netclaw.Actors.Tests/Channels/DiscordGatewayActorTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@
using Microsoft.Extensions.Hosting;
using Netclaw.Actors.Channels;
using Netclaw.Actors.Protocol;
using Netclaw.Actors.Reminders;
using Netclaw.Actors.Tests.Channels.TestHelpers;
using Netclaw.Channels.Discord;
using Netclaw.Configuration;
Expand Down Expand Up @@ -251,7 +252,7 @@ public async Task Gateway_routes_trusted_session_turn_to_conversation_actor()
{
SourceKind = new Netclaw.Actors.Channels.SourceKind("reminder")
},
ReminderId = "rem-1"
ReminderId = new ReminderId("rem-1")
});

gateway.Tell(turn);
Expand Down Expand Up @@ -285,7 +286,7 @@ public async Task Gateway_nacks_trusted_session_turn_with_invalid_session_id()
{
SourceKind = new Netclaw.Actors.Channels.SourceKind("reminder")
},
ReminderId = "rem-1"
ReminderId = new ReminderId("rem-1")
});

gateway.Tell(turn, TestActor);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@
using Microsoft.Extensions.Hosting;
using Netclaw.Actors.Channels;
using Netclaw.Actors.Protocol;
using Netclaw.Actors.Reminders;
using Netclaw.Actors.Tests.Channels.TestHelpers;
using Netclaw.Channels.Mattermost;
using Netclaw.Configuration;
Expand Down Expand Up @@ -573,6 +574,6 @@ private static MattermostGatewayMessage CreateMessage(
{
SourceKind = new SourceKind("reminder")
},
ReminderId = "rem-1"
ReminderId = new ReminderId("rem-1")
};
}
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
using Akka.Hosting.TestKit;
using Microsoft.Extensions.AI;
using Netclaw.Actors.Channels;
using Netclaw.Actors.Reminders;
using Netclaw.Configuration;
using Netclaw.Tools;
using Xunit;
Expand Down Expand Up @@ -41,7 +42,7 @@ private static ChannelInput BuildInput(
?? new SourceProvenance(TransportAuthenticity.Verified, PayloadTaint.Public),
Contents = [new TextContent("hello")],
ReceivedAt = DateTimeOffset.UtcNow,
ReminderId = reminderId,
ReminderId = reminderId is null ? null : new ReminderId(reminderId),
AckTarget = ackTarget,
DefaultDeliveryTarget = defaultDeliveryTarget,
RequestedDeliveryTarget = requestedDeliveryTarget,
Expand Down Expand Up @@ -96,7 +97,7 @@ public void Create_propagates_ReminderId_and_AckTarget_from_ChannelInput()
var result = MessageSourceFactory.Create(
input, new SessionPipelineOptions { ChannelType = ChannelType.Slack }, new Netclaw.Actors.Protocol.TurnId("turn-1"));

Assert.Equal("check-pr:1712000000000", result.ReminderId);
Assert.Equal(new ReminderId("check-pr:1712000000000"), result.ReminderId);
Assert.Same(probe.Ref, result.AckTarget);
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@
using Microsoft.Extensions.Hosting;
using Netclaw.Actors.Channels;
using Netclaw.Actors.Protocol;
using Netclaw.Actors.Reminders;
using Netclaw.Actors.Tests.Channels.TestHelpers;
using Netclaw.Channels.Slack;
using Netclaw.Configuration;
Expand Down Expand Up @@ -442,7 +443,7 @@ public async Task Conversation_rejects_DeliverTrustedSessionTurn_for_other_chann
SourceKind = new Netclaw.Actors.Channels.SourceKind("reminder")
},
ReceivedAt = DateTimeOffset.UtcNow,
ReminderId = reminderId
ReminderId = new ReminderId(reminderId)
};

}
6 changes: 4 additions & 2 deletions src/Netclaw.Actors.Tests/Channels/TurnContextTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,9 @@
// -----------------------------------------------------------------------
using Akka.Actor;
using Netclaw.Actors.Channels;
using Netclaw.Actors.Jobs;
using Netclaw.Actors.Protocol;
using Netclaw.Actors.Reminders;
using Netclaw.Actors.Sessions;
using Netclaw.Configuration;
using Netclaw.Tools;
Expand Down Expand Up @@ -54,8 +56,8 @@ public void FromMessageSource_captures_durable_authority_fields()
AdoptedContextProjection = "quoted context",
AdoptedContextLowerBound = "1700000000.000000",
AdoptedContextUpperBound = "1700000000.000001",
ReminderId = "reminder:1700000000000",
BackgroundJobId = "bg-job:42",
ReminderId = new ReminderId("reminder:1700000000000"),
BackgroundJobId = new BackgroundJobId("bg-job:42"),
AckTarget = ActorRefs.Nobody
};

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -98,7 +98,7 @@ public async Task BackgroundJob_Completes_And_DeliversResult_ViaGateway()
Assert.Equal(PrincipalClassification.VerifiedAutomation, delivered.Source.Principal);
Assert.Equal("background-job", delivered.Source.Provenance.SourceKind?.Value);
Assert.NotNull(delivered.Source.BackgroundJobId);
Assert.StartsWith("bg-job:", delivered.Source.BackgroundJobId);
Assert.StartsWith("bg-job:", delivered.Source.BackgroundJobId!.Value.Value);

await AwaitAssertAsync(() =>
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -253,7 +253,7 @@ await manager.Ask<BackgroundJobManagerHealthResponse>(
Assert.Contains(logPath, delivery.Content);
Assert.Contains("Server running on", delivery.Content);
Assert.Equal(TrustAudience.Personal, delivery.Source.Audience);
Assert.Equal($"bg-job:{orphanId.Value}", delivery.Source.BackgroundJobId);
Assert.Equal(new BackgroundJobId($"bg-job:{orphanId.Value}"), delivery.Source.BackgroundJobId);
}

[Fact]
Expand Down
67 changes: 67 additions & 0 deletions src/Netclaw.Actors.Tests/Protocol/SerializationRoundTripTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -9,8 +9,10 @@
using Google.Protobuf;
using Netclaw.Actors.Channels;
using Netclaw.Actors.Hosting;
using Netclaw.Actors.Jobs;
using Netclaw.Actors.Protocol;
using Netclaw.Actors.Reminders;
using Netclaw.Actors.Serialization;
using Netclaw.Actors.Sessions;
using Netclaw.Tools;
using Xunit;
Expand Down Expand Up @@ -130,6 +132,71 @@ public void TurnRecorded_round_trips()
Assert.Equal(original.RecordedAtMs, result.RecordedAtMs);
}

[Fact]
public void TurnRecorded_round_trips_preserving_value_object_source_ids()
{
var original = new TurnRecorded
{
SessionId = new SessionId("C99999/1708531200.000100"),
UserMessage = new SerializableChatMessage { Role = ChatRole.User, Content = "check PR" },
AssistantReply = new SerializableChatMessage { Role = ChatRole.Assistant, Content = "merged" },
RecordedAtMs = 1_700_000_000_000,
SourceReminderId = new ReminderId("check-pr:1712000000000"),
SourceBackgroundJobId = new BackgroundJobId("bg-job:abc123")
};

var result = RoundTrip(original);

Assert.Equal(original.SourceReminderId, result.SourceReminderId);
Assert.Equal(original.SourceBackgroundJobId, result.SourceBackgroundJobId);
}

[Fact]
public void TurnRecorded_round_trips_with_null_source_ids()
{
var original = new TurnRecorded
{
SessionId = new SessionId("C99999/1708531200.000100"),
UserMessage = new SerializableChatMessage { Role = ChatRole.User, Content = "hi" },
AssistantReply = new SerializableChatMessage { Role = ChatRole.Assistant, Content = "hello" },
RecordedAtMs = 1_700_000_000_000
};

var result = RoundTrip(original);

Assert.Null(result.SourceReminderId);
Assert.Null(result.SourceBackgroundJobId);
}

[Fact]
public void TurnRecorded_value_object_source_ids_use_bare_string_proto_fields()
{
// Wire-compat: the reminder/background-job value objects map to the SAME
// bare-string proto fields the pre-value-object code used, so old journals
// deserialize unchanged and new journals are byte-identical. Proven at the
// proto-mapper boundary in both directions.
var evt = new TurnRecorded
{
SessionId = new SessionId("C99999/1708531200.000100"),
UserMessage = new SerializableChatMessage { Role = ChatRole.User, Content = "x" },
AssistantReply = new SerializableChatMessage { Role = ChatRole.Assistant, Content = "y" },
RecordedAtMs = 1_700_000_000_000,
SourceReminderId = new ReminderId("check-pr:1712000000000"),
SourceBackgroundJobId = new BackgroundJobId("bg-job:abc123")
};

// Forward: value object -> bare string on the wire (no nested object).
var proto = NetclawProtoMapper.ToProto(evt);
Assert.Equal("check-pr:1712000000000", proto.SourceReminderId);
Assert.Equal("bg-job:abc123", proto.SourceBackgroundJobId);

// Reverse: an "old" proto carrying bare strings deserializes into the
// value-object-typed event.
var restored = NetclawProtoMapper.FromProto(proto);
Assert.Equal(new ReminderId("check-pr:1712000000000"), restored.SourceReminderId);
Assert.Equal(new BackgroundJobId("bg-job:abc123"), restored.SourceBackgroundJobId);
}

[Fact]
public void SessionCompacted_round_trips_with_messages()
{
Expand Down
10 changes: 5 additions & 5 deletions src/Netclaw.Actors.Tests/Reminders/ReminderManagerActorTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -510,7 +510,7 @@ public async Task Mode_B_reminder_dispatches_to_resolved_gateway_and_completes_o
Assert.Equal(TrustAudience.Team, delivered.Source.Audience);
Assert.Equal(TrustBoundary.TrustedInstance, delivered.Source.Boundary);
Assert.NotNull(delivered.Source.ReminderId);
Assert.StartsWith("mode-b-anchor:", delivered.Source.ReminderId);
Assert.StartsWith("mode-b-anchor:", delivered.Source.ReminderId!.Value.Value);
Assert.Equal(PrincipalClassification.VerifiedAutomation, delivered.Source.Principal);
Assert.Equal("reminder", delivered.Source.Provenance.SourceKind?.Value);

Expand Down Expand Up @@ -659,7 +659,7 @@ public async Task CurrentSession_delivery_required_fails_fast_on_explicit_delive
// Channel reports the post failed — execution must report failure
// (so Akka.Reminders redelivers) without acking the envelope.
delivered.Source.DeliveryObserver!.Tell(new ReminderDeliveryResult(
delivered.Source.ReminderId!,
delivered.Source.ReminderId!.Value,
ChannelType.Slack,
Delivered: false,
FailureReason: "channel API down"));
Expand Down Expand Up @@ -710,7 +710,7 @@ public async Task CurrentSession_delivery_key_is_built_from_envelope_fire_time()

var delivered = await gatewayProbe.ExpectMsgAsync<DeliverTrustedSessionTurn>(
TimeSpan.FromSeconds(5), cancellationToken: TestContext.Current.CancellationToken);
Assert.Equal($"{definition.Id}:{fireTime.ToUnixTimeMilliseconds()}", delivered.Source.ReminderId);
Assert.Equal(new ReminderId($"{definition.Id}:{fireTime.ToUnixTimeMilliseconds()}"), delivered.Source.ReminderId);
}

[Fact]
Expand Down Expand Up @@ -740,7 +740,7 @@ public async Task CurrentSession_delivery_required_succeeds_when_delivery_is_obs
Assert.NotNull(delivered.Source.DeliveryObserver);

delivered.Source.DeliveryObserver!.Tell(new ReminderDeliveryResult(
delivered.Source.ReminderId!,
delivered.Source.ReminderId!.Value,
ChannelType.Slack,
Delivered: true,
ObservedAtMs: TimeProvider.System.GetUtcNow().ToUnixTimeMilliseconds()));
Expand Down Expand Up @@ -1024,7 +1024,7 @@ public async Task CurrentSession_discord_delivery_required_succeeds_when_deliver
Assert.NotNull(delivered.Source.DeliveryObserver);

delivered.Source.DeliveryObserver!.Tell(new ReminderDeliveryResult(
delivered.Source.ReminderId!,
delivered.Source.ReminderId!.Value,
ChannelType.Discord,
Delivered: true,
ObservedAtMs: TimeProvider.System.GetUtcNow().ToUnixTimeMilliseconds()));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -66,13 +66,13 @@ public void TurnRecorded_WithSourceBackgroundJobId_DedupAndRemovesActive()
UserMessage = new SerializableChatMessage { Role = ChatRole.User, Content = "result" },
AssistantReply = new SerializableChatMessage { Role = ChatRole.Assistant, Content = "ok" },
RecordedAtMs = DateTimeOffset.UtcNow.ToUnixTimeMilliseconds(),
SourceBackgroundJobId = jobKey
SourceBackgroundJobId = new BackgroundJobId(jobKey)
};

state = state.Apply(evt);

Assert.Empty(state.ActiveBackgroundJobs);
Assert.Contains(jobKey, state.ProcessedBackgroundJobIds);
Assert.Contains(new BackgroundJobId(jobKey), state.ProcessedBackgroundJobIds);
}

[Fact]
Expand Down Expand Up @@ -209,7 +209,7 @@ public void Compaction_Preserves_ActiveJobsAndDedupSet()
var processedKey = "bg-job:already-done";
state = state with
{
ProcessedBackgroundJobIds = state.ProcessedBackgroundJobIds.Add(processedKey)
ProcessedBackgroundJobIds = state.ProcessedBackgroundJobIds.Add(new BackgroundJobId(processedKey))
};

var compactedEvt = new SessionCompacted
Expand All @@ -225,6 +225,6 @@ public void Compaction_Preserves_ActiveJobsAndDedupSet()

Assert.Single(state.ActiveBackgroundJobs);
Assert.True(state.ActiveBackgroundJobs.ContainsKey(jobKey));
Assert.Contains(processedKey, state.ProcessedBackgroundJobIds);
Assert.Contains(new BackgroundJobId(processedKey), state.ProcessedBackgroundJobIds);
}
}
Loading
Loading