From 1518d4a4de4f2e975ace64b1a154c00f16038adf Mon Sep 17 00:00:00 2001 From: Einar Date: Thu, 24 Sep 2026 22:41:47 +0200 Subject: [PATCH 1/5] Call the observers contract directly instead of looking it up by reflection clear-quarantine has never worked against a real kernel. It resolved ClearObserverQuarantine with Type.GetMethods() on the client, and the client protobuf-net.Grpc generates implements IObservers explicitly - GetMethods() does not return explicitly implemented members, so the lookup found nothing and the command reported that it could not clear the quarantine. Calling the interface method directly lets the compiler bind it. The specifications that covered this passed throughout, which is why it shipped: they ran against a substitute, and Castle's proxy implements the interface implicitly, so reflection found the method there and only there. The double disagreed with the real client on the one property the command depended on. They are replaced with specifications that run against a double implementing IObservers explicitly, like the generated client does. Restoring the reflection lookup now turns those specifications red. Fixes #185 --- .../and_the_kernel_accepts_it.cs | 62 ++++++++++++ .../and_method_does_not_exist.cs | 15 --- .../and_method_exists.cs | 46 --------- ..._client_shaped_like_the_generated_proxy.cs | 59 ++++++++++++ .../ClearObserverQuarantineCommand.cs | 96 ++++--------------- 5 files changed, 137 insertions(+), 141 deletions(-) create mode 100644 Source/Cli.Specs/for_ClearObserverQuarantineCommand/when_clearing_quarantine/and_the_kernel_accepts_it.cs delete mode 100644 Source/Cli.Specs/for_ClearObserverQuarantineCommand/when_invoking_clear_quarantine/and_method_does_not_exist.cs delete mode 100644 Source/Cli.Specs/for_ClearObserverQuarantineCommand/when_invoking_clear_quarantine/and_method_exists.cs create mode 100644 Source/Cli.Specs/given/an_observers_client_shaped_like_the_generated_proxy.cs diff --git a/Source/Cli.Specs/for_ClearObserverQuarantineCommand/when_clearing_quarantine/and_the_kernel_accepts_it.cs b/Source/Cli.Specs/for_ClearObserverQuarantineCommand/when_clearing_quarantine/and_the_kernel_accepts_it.cs new file mode 100644 index 00000000..e5c11b1d --- /dev/null +++ b/Source/Cli.Specs/for_ClearObserverQuarantineCommand/when_clearing_quarantine/and_the_kernel_accepts_it.cs @@ -0,0 +1,62 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. + +using Cratis.Chronicle.Contracts; +using Cratis.Chronicle.Contracts.Observation; +using Cratis.Cli.Commands.Chronicle.Observers; +using ClearObserverQuarantineContract = Cratis.Chronicle.Contracts.Observation.ClearObserverQuarantine; + +namespace Cratis.Cli.for_ClearObserverQuarantineCommand.when_clearing_quarantine; + +/// +/// This replaces a specification that passed for the whole life of issue #185, while the command failed against +/// every kernel it was pointed at. The command resolved its contract method by reflection, which cannot see +/// explicitly implemented members, and the shipped client implements the contract explicitly - but the double it +/// was specified against implemented it implicitly, so the lookup succeeded in the specification and nowhere else. +/// Running against is what closes that gap: +/// the double now fails the way the real client does. +/// +public class and_the_kernel_accepts_it : Specification +{ + IObservers _observers; + IServices _services; + ObserverCommandSettings _settings; + int _result; + ClearObserverQuarantineContract _sent; + + void Establish() + { + _observers = Substitute.For(); + _services = Substitute.For(); + _services.Observers.Returns(new given.an_observers_client_shaped_like_the_generated_proxy(_observers)); + + _settings = new ObserverCommandSettings + { + EventStore = "the-event-store", + Namespace = "the-namespace", + ObserverId = "the-observer", + EventSequenceId = "event-log" + }; + } + + async Task Because() + { + _result = await new ClearObserverQuarantineCommandForSpecs().Execute(_services, _settings); + _sent = _observers.ReceivedCalls() + .Select(call => call.GetArguments()[0]) + .OfType() + .Single(); + } + + [Fact] void should_succeed() => _result.ShouldEqual(ExitCodes.Success); + [Fact] void should_clear_the_observer_it_was_given() => _sent.ObserverId.ShouldEqual("the-observer"); + [Fact] void should_target_the_resolved_event_store() => _sent.EventStore.ShouldEqual("the-event-store"); + [Fact] void should_target_the_resolved_namespace() => _sent.Namespace.ShouldEqual("the-namespace"); + [Fact] void should_pass_the_event_sequence() => _sent.EventSequenceId.ShouldEqual("event-log"); + + class ClearObserverQuarantineCommandForSpecs : ClearObserverQuarantineCommand + { + public Task Execute(IServices services, ObserverCommandSettings settings) => + ExecuteCommandAsync(services, settings, "json"); + } +} diff --git a/Source/Cli.Specs/for_ClearObserverQuarantineCommand/when_invoking_clear_quarantine/and_method_does_not_exist.cs b/Source/Cli.Specs/for_ClearObserverQuarantineCommand/when_invoking_clear_quarantine/and_method_does_not_exist.cs deleted file mode 100644 index 1eefc0e0..00000000 --- a/Source/Cli.Specs/for_ClearObserverQuarantineCommand/when_invoking_clear_quarantine/and_method_does_not_exist.cs +++ /dev/null @@ -1,15 +0,0 @@ -// Copyright (c) Cratis. All rights reserved. -// Licensed under the MIT license. See LICENSE file in the project root for full license information. - -namespace Cratis.Cli.for_ClearObserverQuarantineCommand.when_invoking_clear_quarantine; - -public class and_method_does_not_exist : Specification -{ - bool _result; - - async Task Because() => _result = await ClearObserverQuarantineInvoker.TryClear(new FakeObservers(), "event-store", "namespace", "observer-id", "event-log"); - - [Fact] void should_return_false() => _result.ShouldBeFalse(); - - class FakeObservers; -} diff --git a/Source/Cli.Specs/for_ClearObserverQuarantineCommand/when_invoking_clear_quarantine/and_method_exists.cs b/Source/Cli.Specs/for_ClearObserverQuarantineCommand/when_invoking_clear_quarantine/and_method_exists.cs deleted file mode 100644 index 0515ebd1..00000000 --- a/Source/Cli.Specs/for_ClearObserverQuarantineCommand/when_invoking_clear_quarantine/and_method_exists.cs +++ /dev/null @@ -1,46 +0,0 @@ -// Copyright (c) Cratis. All rights reserved. -// Licensed under the MIT license. See LICENSE file in the project root for full license information. - -namespace Cratis.Cli.for_ClearObserverQuarantineCommand.when_invoking_clear_quarantine; - -public class and_method_exists : Specification -{ - FakeObservers _observers; - bool _result; - - void Establish() => _observers = new(); - - async Task Because() => _result = await ClearObserverQuarantineInvoker.TryClear(_observers, "event-store", "namespace", "observer-id", "event-log"); - - [Fact] void should_invoke_clear_quarantine() => _observers.WasInvoked.ShouldBeTrue(); - [Fact] void should_return_true() => _result.ShouldBeTrue(); - [Fact] void should_set_event_store() => _observers.Command.EventStore.ShouldEqual("event-store"); - [Fact] void should_set_namespace() => _observers.Command.Namespace.ShouldEqual("namespace"); - [Fact] void should_set_observer_id() => _observers.Command.ObserverId.ShouldEqual("observer-id"); - [Fact] void should_set_event_sequence_id() => _observers.Command.EventSequenceId.ShouldEqual("event-log"); - - class FakeObservers - { - public bool WasInvoked { get; private set; } - - public FakeClearObserverQuarantine Command { get; private set; } = null!; - - public Task ClearObserverQuarantine(FakeClearObserverQuarantine command, int ignored = 0) - { - WasInvoked = true; - Command = command; - return Task.CompletedTask; - } - } - - class FakeClearObserverQuarantine - { - public string EventStore { get; set; } = string.Empty; - - public string Namespace { get; set; } = string.Empty; - - public string ObserverId { get; set; } = string.Empty; - - public string EventSequenceId { get; set; } = string.Empty; - } -} diff --git a/Source/Cli.Specs/given/an_observers_client_shaped_like_the_generated_proxy.cs b/Source/Cli.Specs/given/an_observers_client_shaped_like_the_generated_proxy.cs new file mode 100644 index 00000000..5498d750 --- /dev/null +++ b/Source/Cli.Specs/given/an_observers_client_shaped_like_the_generated_proxy.cs @@ -0,0 +1,59 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. + +using Cratis.Chronicle.Contracts.Clients; +using Cratis.Chronicle.Contracts.Observation; +using ProtoBuf.Grpc; + +namespace Cratis.Cli.given; + +/// +/// An that implements every member explicitly, forwarding to a substitute. +/// +/// +/// +/// This exists because of issue #185, and because the specification that was supposed to cover it could not fail. +/// The command looked its contract method up with Type.GetMethods(), which does not return explicitly +/// implemented members - and the real client, generated by protobuf-net.Grpc, implements +/// explicitly. So the lookup found nothing and clear-quarantine failed against every kernel +/// it was ever pointed at. +/// +/// +/// A plain Substitute.For<IObservers>() would not have caught it: Castle's proxy implements the +/// interface implicitly, so reflection finds the method and the specification goes green while production +/// stays broken. That mismatch between the double and the real client is the whole defect, so specifications that +/// care about how the command reaches the contract use this instead - it is shaped like what actually ships. +/// +/// +/// The substitute every call is forwarded to, for stubbing and verification. +public sealed class an_observers_client_shaped_like_the_generated_proxy(IObservers inner) : IObservers +{ + /// + /// Gets the substitute calls are forwarded to. Stub and assert against this. + /// + public IObservers Inner { get; } = inner; + + Task IObservers.Replay(Replay command, CallContext context) => Inner.Replay(command, context); + + Task IObservers.ReplayPartition(ReplayPartition command, CallContext context) => Inner.ReplayPartition(command, context); + + Task IObservers.RetryPartition(RetryPartition command, CallContext context) => Inner.RetryPartition(command, context); + + Task IObservers.ClearObserverQuarantine(ClearObserverQuarantine command, CallContext context) => Inner.ClearObserverQuarantine(command, context); + + Task IObservers.ClearFailedPartitions(ClearFailedPartitions command, CallContext context) => Inner.ClearFailedPartitions(command, context); + + Task IObservers.RemoveObserver(RemoveObserver command, CallContext context) => Inner.RemoveObserver(command, context); + + Task IObservers.GetObserverInformation(GetObserverInformationRequest request, CallContext context) => Inner.GetObserverInformation(request, context); + + Task> IObservers.GetConnectedClientsForObserver(GetConnectedClientsForObserverRequest request, CallContext context) => Inner.GetConnectedClientsForObserver(request, context); + + Task> IObservers.GetObservers(AllObserversRequest request, CallContext context) => Inner.GetObservers(request, context); + + IObservable> IObservers.ObserveObservers(AllObserversRequest request, CallContext context) => Inner.ObserveObservers(request, context); + + Task IObservers.WaitForCompletion(WaitForObserverCompletionRequest request, CallContext context) => Inner.WaitForCompletion(request, context); + + Task> IObservers.GetReplayableObserversForEventTypes(GetReplayableObserversForEventTypesRequest request, CallContext context) => Inner.GetReplayableObserversForEventTypes(request, context); +} diff --git a/Source/Cli/Commands/Chronicle/Observers/ClearObserverQuarantineCommand.cs b/Source/Cli/Commands/Chronicle/Observers/ClearObserverQuarantineCommand.cs index b11621be..8cd7b04a 100644 --- a/Source/Cli/Commands/Chronicle/Observers/ClearObserverQuarantineCommand.cs +++ b/Source/Cli/Commands/Chronicle/Observers/ClearObserverQuarantineCommand.cs @@ -1,11 +1,21 @@ // Copyright (c) Cratis. All rights reserved. // Licensed under the MIT license. See LICENSE file in the project root for full license information. +using ClearObserverQuarantineContract = Cratis.Chronicle.Contracts.Observation.ClearObserverQuarantine; + namespace Cratis.Cli.Commands.Chronicle.Observers; /// /// Clears quarantine for an observer. /// +/// +/// This used to look the contract method up by reflection and report every miss as "the connected kernel does not +/// support this", which is how it came to fail against every kernel that does: the gRPC client is a generated proxy +/// implementing the interface explicitly, and does not return explicit interface +/// implementations. The suggestion it printed - upgrade the contracts - was unfollowable, because there was no version +/// where the lookup would have succeeded. The CLI compiles against the contracts that declare the method, so calling +/// it directly is both simpler and the only way to tell a kernel that cannot do this from a lookup that went wrong. +/// [LlmDescription("Clears quarantine for a quarantined observer so it can resume processing. Prompts for confirmation unless --yes is specified.")] [CommandEffect(CommandEffect.Mutating)] [CliCommand("clear-quarantine", "Clear quarantine for an observer", Branch = typeof(ChronicleBranch.Observers), DynamicCompletion = "observers")] @@ -20,89 +30,15 @@ protected override string GetConfirmationPrompt(ObserverCommandSettings settings /// protected override async Task ExecuteCommandAsync(IServices services, ObserverCommandSettings settings, string format) { - var cleared = await ClearObserverQuarantineInvoker.TryClear( - services.Observers, - settings.ResolveEventStore(), - settings.ResolveNamespace(), - settings.ObserverId, - settings.EventSequenceId); - - if (!cleared) + await services.Observers.ClearObserverQuarantine(new ClearObserverQuarantineContract { - OutputFormatter.WriteError( - format, - "Connected Chronicle contracts do not support clearing observer quarantine.", - "Upgrade Cratis.Chronicle.Contracts and Cratis.Chronicle.Connections to a version that includes observer quarantine clear support.", - ExitCodes.ValidationErrorCode); - return ExitCodes.ValidationError; - } + EventStore = settings.ResolveEventStore(), + Namespace = settings.ResolveNamespace(), + ObserverId = settings.ObserverId, + EventSequenceId = settings.EventSequenceId + }); OutputFormatter.WriteMessage(format, $"Quarantine cleared for observer '{settings.ObserverId}'."); return ExitCodes.Success; } } - -static class ClearObserverQuarantineInvoker -{ - public static async Task TryClear(object observers, string eventStore, string @namespace, string observerId, string eventSequenceId) - { - var method = observers.GetType() - .GetMethods() - .FirstOrDefault(_ => _.Name == "ClearObserverQuarantine" && _.GetParameters().Length > 0); - - if (method is null) - { - return false; - } - - var parameters = method.GetParameters(); - var command = Activator.CreateInstance(parameters[0].ParameterType); - if (command is null) - { - return false; - } - - if (!SetPropertyIfPresent(command, "EventStore", eventStore) || - !SetPropertyIfPresent(command, "Namespace", @namespace) || - !SetPropertyIfPresent(command, "ObserverId", observerId)) - { - return false; - } - - SetPropertyIfPresent(command, "EventSequenceId", eventSequenceId); - - var arguments = new object?[parameters.Length]; - arguments[0] = command; - for (var i = 1; i < parameters.Length; i++) - { - if (parameters[i].HasDefaultValue) - { - arguments[i] = parameters[i].DefaultValue; - } - else if (parameters[i].ParameterType.IsValueType) - { - arguments[i] = Activator.CreateInstance(parameters[i].ParameterType); - } - } - - if (method.Invoke(observers, arguments) is not Task task) - { - return false; - } - - await task; - return true; - } - - static bool SetPropertyIfPresent(object instance, string propertyName, string value) - { - var property = instance.GetType().GetProperty(propertyName); - if (property?.CanWrite != true) - { - return false; - } - - property.SetValue(instance, value); - return true; - } -} From f4c2bcb1ef2f886830d37a81d5b6ef404309af9c Mon Sep 17 00:00:00 2001 From: Einar Date: Thu, 24 Sep 2026 22:41:52 +0200 Subject: [PATCH 2/5] Report what the kernel did with a partition retry instead of assuming it started retry-partition awaited the call and then printed that the retry had started, whatever came back. The kernel declines outright in three cases - the observer is quarantined, so recovery is paused by design; the partition is individually quarantined; the partition is not among the observer's failures - and in each of them the command told the operator to go and watch progress that was never coming. The kernel already reports which of those it did on RetryPartitionResponse, so the command now reads the outcome, says what happened, and exits non-zero when nothing was started. A script can no longer take silence for recovery. Fixes #186 --- .../given/a_retry_partition_command.cs | 52 +++++++++++++++++++ .../and_the_observer_is_quarantined.cs | 23 ++++++++ .../and_the_partition_is_not_failed.cs | 17 ++++++ .../and_the_partition_is_quarantined.cs | 17 ++++++ .../when_retrying/and_the_retry_starts.cs | 13 +++++ .../Observers/RetryPartitionCommand.cs | 43 +++++++++++++-- 6 files changed, 161 insertions(+), 4 deletions(-) create mode 100644 Source/Cli.Specs/for_RetryPartitionCommand/given/a_retry_partition_command.cs create mode 100644 Source/Cli.Specs/for_RetryPartitionCommand/when_retrying/and_the_observer_is_quarantined.cs create mode 100644 Source/Cli.Specs/for_RetryPartitionCommand/when_retrying/and_the_partition_is_not_failed.cs create mode 100644 Source/Cli.Specs/for_RetryPartitionCommand/when_retrying/and_the_partition_is_quarantined.cs create mode 100644 Source/Cli.Specs/for_RetryPartitionCommand/when_retrying/and_the_retry_starts.cs diff --git a/Source/Cli.Specs/for_RetryPartitionCommand/given/a_retry_partition_command.cs b/Source/Cli.Specs/for_RetryPartitionCommand/given/a_retry_partition_command.cs new file mode 100644 index 00000000..14d2a3a7 --- /dev/null +++ b/Source/Cli.Specs/for_RetryPartitionCommand/given/a_retry_partition_command.cs @@ -0,0 +1,52 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. + +using Cratis.Chronicle.Contracts; +using Cratis.Chronicle.Contracts.Observation; +using Cratis.Cli.Commands.Chronicle.Observers; +using Cratis.Cli.given; +using RetryPartitionContract = Cratis.Chronicle.Contracts.Observation.RetryPartition; + +namespace Cratis.Cli.for_RetryPartitionCommand.given; + +public class a_retry_partition_command : Specification +{ + protected IServices _services; + protected IObservers _observers; + protected PartitionCommandSettings _settings; + protected RetryPartitionCommandForSpecs _command; + + void Establish() + { + _observers = Substitute.For(); + _services = Substitute.For(); + + // Explicitly implemented, like the shipped client - see the double's own remarks and issue #185. + _services.Observers.Returns(new an_observers_client_shaped_like_the_generated_proxy(_observers)); + + _settings = new PartitionCommandSettings + { + EventStore = "the-event-store", + Namespace = "the-namespace", + ObserverId = "the-observer", + Partition = "the-partition", + EventSequenceId = "event-log" + }; + + _command = new RetryPartitionCommandForSpecs(); + Recovery(PartitionRecoveryOutcome.Started); + } + + protected void Recovery(PartitionRecoveryOutcome outcome) => + _observers + .RetryPartition(Arg.Any(), Arg.Any()) + .Returns(new RetryPartitionResponse { Outcome = outcome }); + + protected Task Execute() => _command.Execute(_services, _settings); + + public class RetryPartitionCommandForSpecs : RetryPartitionCommand + { + public Task Execute(IServices services, PartitionCommandSettings settings) => + ExecuteCommandAsync(services, settings, "json"); + } +} diff --git a/Source/Cli.Specs/for_RetryPartitionCommand/when_retrying/and_the_observer_is_quarantined.cs b/Source/Cli.Specs/for_RetryPartitionCommand/when_retrying/and_the_observer_is_quarantined.cs new file mode 100644 index 00000000..ab789e3d --- /dev/null +++ b/Source/Cli.Specs/for_RetryPartitionCommand/when_retrying/and_the_observer_is_quarantined.cs @@ -0,0 +1,23 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. + +using Cratis.Chronicle.Contracts.Observation; + +namespace Cratis.Cli.for_RetryPartitionCommand.when_retrying; + +/// +/// A quarantined observer has recovery paused by design, so the retry goes nowhere. This used to print "Retry started" +/// and exit zero, sending the operator off to watch progress that was never coming - which is how a production +/// recovery came down to hand-editing MongoDB. +/// +public class and_the_observer_is_quarantined : given.a_retry_partition_command +{ + int _result; + + void Establish() => Recovery(PartitionRecoveryOutcome.ObserverQuarantined); + + async Task Because() => _result = await Execute(); + + [Fact] void should_not_report_success() => _result.ShouldNotEqual(ExitCodes.Success); + [Fact] void should_report_a_validation_error() => _result.ShouldEqual(ExitCodes.ValidationError); +} diff --git a/Source/Cli.Specs/for_RetryPartitionCommand/when_retrying/and_the_partition_is_not_failed.cs b/Source/Cli.Specs/for_RetryPartitionCommand/when_retrying/and_the_partition_is_not_failed.cs new file mode 100644 index 00000000..c839db37 --- /dev/null +++ b/Source/Cli.Specs/for_RetryPartitionCommand/when_retrying/and_the_partition_is_not_failed.cs @@ -0,0 +1,17 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. + +using Cratis.Chronicle.Contracts.Observation; + +namespace Cratis.Cli.for_RetryPartitionCommand.when_retrying; + +public class and_the_partition_is_not_failed : given.a_retry_partition_command +{ + int _result; + + void Establish() => Recovery(PartitionRecoveryOutcome.PartitionNotFound); + + async Task Because() => _result = await Execute(); + + [Fact] void should_not_report_success() => _result.ShouldEqual(ExitCodes.ValidationError); +} diff --git a/Source/Cli.Specs/for_RetryPartitionCommand/when_retrying/and_the_partition_is_quarantined.cs b/Source/Cli.Specs/for_RetryPartitionCommand/when_retrying/and_the_partition_is_quarantined.cs new file mode 100644 index 00000000..bdeb86a1 --- /dev/null +++ b/Source/Cli.Specs/for_RetryPartitionCommand/when_retrying/and_the_partition_is_quarantined.cs @@ -0,0 +1,17 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. + +using Cratis.Chronicle.Contracts.Observation; + +namespace Cratis.Cli.for_RetryPartitionCommand.when_retrying; + +public class and_the_partition_is_quarantined : given.a_retry_partition_command +{ + int _result; + + void Establish() => Recovery(PartitionRecoveryOutcome.PartitionQuarantined); + + async Task Because() => _result = await Execute(); + + [Fact] void should_not_report_success() => _result.ShouldEqual(ExitCodes.ValidationError); +} diff --git a/Source/Cli.Specs/for_RetryPartitionCommand/when_retrying/and_the_retry_starts.cs b/Source/Cli.Specs/for_RetryPartitionCommand/when_retrying/and_the_retry_starts.cs new file mode 100644 index 00000000..71845c06 --- /dev/null +++ b/Source/Cli.Specs/for_RetryPartitionCommand/when_retrying/and_the_retry_starts.cs @@ -0,0 +1,13 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. + +namespace Cratis.Cli.for_RetryPartitionCommand.when_retrying; + +public class and_the_retry_starts : given.a_retry_partition_command +{ + int _result; + + async Task Because() => _result = await Execute(); + + [Fact] void should_succeed() => _result.ShouldEqual(ExitCodes.Success); +} diff --git a/Source/Cli/Commands/Chronicle/Observers/RetryPartitionCommand.cs b/Source/Cli/Commands/Chronicle/Observers/RetryPartitionCommand.cs index fc84e4c6..f46c9b30 100644 --- a/Source/Cli/Commands/Chronicle/Observers/RetryPartitionCommand.cs +++ b/Source/Cli/Commands/Chronicle/Observers/RetryPartitionCommand.cs @@ -1,6 +1,7 @@ // Copyright (c) Cratis. All rights reserved. // Licensed under the MIT license. See LICENSE file in the project root for full license information. +using Cratis.Chronicle.Contracts.Observation; using RetryPartitionContract = Cratis.Chronicle.Contracts.Observation.RetryPartition; namespace Cratis.Cli.Commands.Chronicle.Observers; @@ -8,7 +9,14 @@ namespace Cratis.Cli.Commands.Chronicle.Observers; /// /// Retries a failed partition. /// -[LlmDescription("Retries a failed partition of an observer from the sequence number where it last failed. Use to recover a failed partition after the underlying issue is resolved.")] +/// +/// This used to print "Retry started" and exit zero whatever happened, including for the cases the kernel declines +/// outright - a quarantined observer has recovery paused by design, an individually quarantined partition is not +/// retried until it is cleared, and a partition that is not among the observer's failures has nothing to recover. An +/// operator recovering a stuck store was told to go watch progress that was never coming. The kernel reports which of +/// those it did, so the command says so and exits non-zero when nothing was started. +/// +[LlmDescription("Retries a failed partition of an observer from the sequence number where it last failed. Use to recover a failed partition after the underlying issue is resolved. Exits non-zero when the retry is refused.")] [CommandEffect(CommandEffect.Mutating)] [CliCommand("retry-partition", "Retry a failed partition", Branch = typeof(ChronicleBranch.Observers), DynamicCompletion = "observers")] [CliExample("chronicle", "observers", "retry-partition", "550e8400-e29b-41d4-a716-446655440000", "my-partition")] @@ -23,7 +31,7 @@ protected override string GetConfirmationPrompt(PartitionCommandSettings setting /// protected override async Task ExecuteCommandAsync(IServices services, PartitionCommandSettings settings, string format) { - await services.Observers.RetryPartition(new RetryPartitionContract + var response = await services.Observers.RetryPartition(new RetryPartitionContract { EventStore = settings.ResolveEventStore(), Namespace = settings.ResolveNamespace(), @@ -32,7 +40,34 @@ await services.Observers.RetryPartition(new RetryPartitionContract Partition = settings.Partition }); - OutputFormatter.WriteMessage(format, $"Retry started for partition '{settings.Partition}' of observer '{settings.ObserverId}'. Use 'cratis observers show {settings.ObserverId}' to check progress."); - return ExitCodes.Success; + if (response.Outcome == PartitionRecoveryOutcome.Started) + { + OutputFormatter.WriteMessage(format, $"Retry started for partition '{settings.Partition}' of observer '{settings.ObserverId}'. Use 'cratis chronicle observers show {settings.ObserverId}' to check progress."); + return ExitCodes.Success; + } + + var (message, suggestion) = DescribeRefusal(response.Outcome, settings); + OutputFormatter.WriteError(format, message, suggestion, ExitCodes.ValidationErrorCode); + return ExitCodes.ValidationError; } + + static (string Message, string Suggestion) DescribeRefusal(PartitionRecoveryOutcome outcome, PartitionCommandSettings settings) => + outcome switch + { + PartitionRecoveryOutcome.ObserverQuarantined => + ($"Observer '{settings.ObserverId}' is quarantined, so recovery is paused and partition '{settings.Partition}' was not retried.", + $"Clear the observer quarantine first: cratis chronicle observers clear-quarantine {settings.ObserverId}"), + + PartitionRecoveryOutcome.PartitionQuarantined => + ($"Partition '{settings.Partition}' of observer '{settings.ObserverId}' is quarantined after exhausting its retry attempts, so it was not retried.", + "Clear the partition quarantine before retrying it."), + + PartitionRecoveryOutcome.PartitionNotFound => + ($"Partition '{settings.Partition}' is not among the failed partitions of observer '{settings.ObserverId}', so there was nothing to retry.", + $"List the failed partitions to check the key: cratis chronicle failed-partitions list --observer {settings.ObserverId}"), + + _ => + ($"Partition '{settings.Partition}' of observer '{settings.ObserverId}' was not retried: {outcome}.", + "Check the observer state with 'cratis chronicle observers show'.") + }; } From b14c859ef9a14eb282791b7f03f4e2014aa1181e Mon Sep 17 00:00:00 2001 From: Einar Date: Thu, 24 Sep 2026 22:41:58 +0200 Subject: [PATCH 3/5] Add observers remove for observers whose declaring code is gone A deleted read model and its projection, or a removed reactor, leaves its definition, state, handled counts and failed partitions in the event store with nothing left to claim them and no supported way to clear them. The command asks first, and says what it is about to do: what is deleted, that it reaches every namespace of the event store, that read models and their data are left alone, that it cannot be undone, and that a re-registered observer starts over from the beginning of the sequence - which is the distinction from a replay, and what an operator reaching for this may actually have meant. The kernel refuses while the observer is running or a client is still subscribed to it, so the command reports the refusal and exits non-zero rather than leaving the impression that the store was cleaned up. --- .../given/a_remove_observer_command.cs | 56 +++++++++++ .../and_the_operator_is_asked.cs | 23 +++++ .../and_the_observer_does_not_exist.cs | 17 ++++ .../and_the_observer_is_removed.cs | 27 ++++++ .../and_the_observer_is_still_running.cs | 23 +++++ .../Observers/RemoveObserverCommand.cs | 96 +++++++++++++++++++ 6 files changed, 242 insertions(+) create mode 100644 Source/Cli.Specs/for_RemoveObserverCommand/given/a_remove_observer_command.cs create mode 100644 Source/Cli.Specs/for_RemoveObserverCommand/when_confirming/and_the_operator_is_asked.cs create mode 100644 Source/Cli.Specs/for_RemoveObserverCommand/when_removing/and_the_observer_does_not_exist.cs create mode 100644 Source/Cli.Specs/for_RemoveObserverCommand/when_removing/and_the_observer_is_removed.cs create mode 100644 Source/Cli.Specs/for_RemoveObserverCommand/when_removing/and_the_observer_is_still_running.cs create mode 100644 Source/Cli/Commands/Chronicle/Observers/RemoveObserverCommand.cs diff --git a/Source/Cli.Specs/for_RemoveObserverCommand/given/a_remove_observer_command.cs b/Source/Cli.Specs/for_RemoveObserverCommand/given/a_remove_observer_command.cs new file mode 100644 index 00000000..2f90cd73 --- /dev/null +++ b/Source/Cli.Specs/for_RemoveObserverCommand/given/a_remove_observer_command.cs @@ -0,0 +1,56 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. + +using Cratis.Chronicle.Contracts; +using Cratis.Chronicle.Contracts.Observation; +using Cratis.Cli.Commands.Chronicle.Observers; +using Cratis.Cli.given; +using RemoveObserverContract = Cratis.Chronicle.Contracts.Observation.RemoveObserver; + +namespace Cratis.Cli.for_RemoveObserverCommand.given; + +public class a_remove_observer_command : Specification +{ + protected IServices _services; + protected IObservers _observers; + protected ObserverCommandSettings _settings; + protected RemoveObserverCommandForSpecs _command; + + void Establish() + { + _observers = Substitute.For(); + _services = Substitute.For(); + + // Explicitly implemented, like the shipped client - see the double's own remarks and issue #185. + _services.Observers.Returns(new an_observers_client_shaped_like_the_generated_proxy(_observers)); + + _settings = new ObserverCommandSettings + { + EventStore = "the-event-store", + Namespace = "the-namespace", + ObserverId = "the-observer", + EventSequenceId = "event-log" + }; + + _command = new RemoveObserverCommandForSpecs(); + Removal(ObserverRemovalOutcome.Removed); + } + + protected void Removal(ObserverRemovalOutcome outcome, string blockingNamespace = "") => + _observers + .RemoveObserver(Arg.Any(), Arg.Any()) + .Returns(new RemoveObserverResponse { Outcome = outcome, BlockingNamespace = blockingNamespace }); + + protected Task Execute() => _command.Execute(_services, _settings); + + /// + /// Exposes the command's execution and its confirmation prompt, both of which are protected on the base command. + /// + public class RemoveObserverCommandForSpecs : RemoveObserverCommand + { + public Task Execute(IServices services, ObserverCommandSettings settings) => + ExecuteCommandAsync(services, settings, "json"); + + public string Confirmation(ObserverCommandSettings settings) => GetConfirmationPrompt(settings); + } +} diff --git a/Source/Cli.Specs/for_RemoveObserverCommand/when_confirming/and_the_operator_is_asked.cs b/Source/Cli.Specs/for_RemoveObserverCommand/when_confirming/and_the_operator_is_asked.cs new file mode 100644 index 00000000..ea771b49 --- /dev/null +++ b/Source/Cli.Specs/for_RemoveObserverCommand/when_confirming/and_the_operator_is_asked.cs @@ -0,0 +1,23 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. + +namespace Cratis.Cli.for_RemoveObserverCommand.when_confirming; + +/// +/// Removal cannot be undone and reaches every namespace in the event store, so the prompt has to say what goes, what +/// stays and that a re-registered observer starts over - the distinction from a replay, which is what an operator +/// reaching for this may well have meant. A bare "are you sure?" gives no basis for answering. +/// +public class and_the_operator_is_asked : given.a_remove_observer_command +{ + string _prompt; + + void Because() => _prompt = _command.Confirmation(_settings); + + [Fact] void should_name_the_observer() => _prompt.ShouldContain("the-observer"); + [Fact] void should_name_the_event_store() => _prompt.ShouldContain("the-event-store"); + [Fact] void should_say_it_covers_every_namespace() => _prompt.ShouldContain("every namespace"); + [Fact] void should_say_read_models_are_left_alone() => _prompt.ShouldContain("Read models and their data are not touched"); + [Fact] void should_say_it_cannot_be_undone() => _prompt.ShouldContain("cannot be undone"); + [Fact] void should_say_a_re_registered_observer_starts_over() => _prompt.ShouldContain("start over"); +} diff --git a/Source/Cli.Specs/for_RemoveObserverCommand/when_removing/and_the_observer_does_not_exist.cs b/Source/Cli.Specs/for_RemoveObserverCommand/when_removing/and_the_observer_does_not_exist.cs new file mode 100644 index 00000000..abf6f6b3 --- /dev/null +++ b/Source/Cli.Specs/for_RemoveObserverCommand/when_removing/and_the_observer_does_not_exist.cs @@ -0,0 +1,17 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. + +using Cratis.Chronicle.Contracts.Observation; + +namespace Cratis.Cli.for_RemoveObserverCommand.when_removing; + +public class and_the_observer_does_not_exist : given.a_remove_observer_command +{ + int _result; + + void Establish() => Removal(ObserverRemovalOutcome.ObserverNotFound); + + async Task Because() => _result = await Execute(); + + [Fact] void should_report_it_as_not_found() => _result.ShouldEqual(ExitCodes.NotFound); +} diff --git a/Source/Cli.Specs/for_RemoveObserverCommand/when_removing/and_the_observer_is_removed.cs b/Source/Cli.Specs/for_RemoveObserverCommand/when_removing/and_the_observer_is_removed.cs new file mode 100644 index 00000000..bbdb2e01 --- /dev/null +++ b/Source/Cli.Specs/for_RemoveObserverCommand/when_removing/and_the_observer_is_removed.cs @@ -0,0 +1,27 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. + +using RemoveObserverContract = Cratis.Chronicle.Contracts.Observation.RemoveObserver; + +namespace Cratis.Cli.for_RemoveObserverCommand.when_removing; + +public class and_the_observer_is_removed : given.a_remove_observer_command +{ + int _result; + RemoveObserverContract _sent; + + async Task Because() + { + _result = await Execute(); + _sent = _observers.ReceivedCalls() + .Select(call => call.GetArguments()[0]) + .OfType() + .Single(); + } + + [Fact] void should_succeed() => _result.ShouldEqual(ExitCodes.Success); + [Fact] void should_remove_the_observer_it_was_given() => _sent.ObserverId.ShouldEqual("the-observer"); + [Fact] void should_target_the_resolved_event_store() => _sent.EventStore.ShouldEqual("the-event-store"); + [Fact] void should_target_the_resolved_namespace() => _sent.Namespace.ShouldEqual("the-namespace"); + [Fact] void should_pass_the_event_sequence() => _sent.EventSequenceId.ShouldEqual("event-log"); +} diff --git a/Source/Cli.Specs/for_RemoveObserverCommand/when_removing/and_the_observer_is_still_running.cs b/Source/Cli.Specs/for_RemoveObserverCommand/when_removing/and_the_observer_is_still_running.cs new file mode 100644 index 00000000..db25667f --- /dev/null +++ b/Source/Cli.Specs/for_RemoveObserverCommand/when_removing/and_the_observer_is_still_running.cs @@ -0,0 +1,23 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. + +using Cratis.Chronicle.Contracts.Observation; + +namespace Cratis.Cli.for_RemoveObserverCommand.when_removing; + +/// +/// A refusal has to exit non-zero. It comes back as a perfectly successful RPC that declined, so reporting the +/// outcome is the only thing that separates "the observer is gone" from "nothing happened" - and a script that only +/// looks at the exit code would otherwise carry on as though the store had been cleaned up. +/// +public class and_the_observer_is_still_running : given.a_remove_observer_command +{ + int _result; + + void Establish() => Removal(ObserverRemovalOutcome.ObserverActive, "the-busy-namespace"); + + async Task Because() => _result = await Execute(); + + [Fact] void should_not_report_success() => _result.ShouldNotEqual(ExitCodes.Success); + [Fact] void should_report_a_validation_error() => _result.ShouldEqual(ExitCodes.ValidationError); +} diff --git a/Source/Cli/Commands/Chronicle/Observers/RemoveObserverCommand.cs b/Source/Cli/Commands/Chronicle/Observers/RemoveObserverCommand.cs new file mode 100644 index 00000000..0e70fe03 --- /dev/null +++ b/Source/Cli/Commands/Chronicle/Observers/RemoveObserverCommand.cs @@ -0,0 +1,96 @@ +// Copyright (c) Cratis. All rights reserved. +// Licensed under the MIT license. See LICENSE file in the project root for full license information. + +using Cratis.Chronicle.Contracts.Observation; +using RemoveObserverContract = Cratis.Chronicle.Contracts.Observation.RemoveObserver; + +namespace Cratis.Cli.Commands.Chronicle.Observers; + +/// +/// Removes an observer and everything the event store keeps for it. +/// +/// +/// For the observer whose declaring code is gone - a read model and its projection that were deleted, a reactor that +/// was removed. Nothing tells the event store, so the observer settles into Disconnected and keeps its records +/// forever. The kernel refuses while the observer is running or still has a subscribed client, so this cannot be used +/// to tear a live observer out from under a running application, and there is deliberately no force flag. +/// +[LlmDescription("Removes an observer and everything the event store keeps for it - its definition, state, handled counts, failed partitions and, for a projection, its projection definition - across every namespace of the event store. Read model data is not touched. Refused while the observer is running or still has a subscribed client. Prompts for confirmation unless --yes is specified.")] +[CommandEffect(CommandEffect.Destructive)] +[CliCommand("remove", "Remove an observer and everything kept for it", Branch = typeof(ChronicleBranch.Observers), DynamicCompletion = "observers")] +[CliExample("chronicle", "observers", "remove", "550e8400-e29b-41d4-a716-446655440000")] +[LlmOption("", "string", "Observer identifier (from 'cratis chronicle observers list') (positional)")] +public class RemoveObserverCommand : ChronicleCommand +{ + /// + /// + /// The prompt says what goes, what stays and that it reaches every namespace. Removal cannot be undone, and the + /// difference between it and a replay - which an operator reaching for it may well have meant - is that a + /// re-registered observer starts over from the beginning of the sequence. + /// + protected override string GetConfirmationPrompt(ObserverCommandSettings settings) => + string.Join( + Environment.NewLine, + [ + $"Remove observer '{settings.ObserverId}' and everything the event store keeps for it?", + string.Empty, + "This deletes its definition, its state and handled counts in every namespace of event store", + $"'{settings.ResolveEventStore()}', its failed partitions and, if it is a projection, its projection", + "definition. Read models and their data are not touched.", + string.Empty, + "This cannot be undone. The observer comes back only if the code that declared it registers it again,", + "and it will then start over from the beginning of the event sequence.", + string.Empty, + "Are you sure?" + ]); + + /// + protected override async Task ExecuteCommandAsync(IServices services, ObserverCommandSettings settings, string format) + { + var response = await services.Observers.RemoveObserver(new RemoveObserverContract + { + EventStore = settings.ResolveEventStore(), + Namespace = settings.ResolveNamespace(), + ObserverId = settings.ObserverId, + EventSequenceId = settings.EventSequenceId + }); + + if (response.Outcome == ObserverRemovalOutcome.Removed) + { + OutputFormatter.WriteMessage(format, $"Observer '{settings.ObserverId}' was removed from event store '{settings.ResolveEventStore()}'."); + return ExitCodes.Success; + } + + var (message, suggestion, exitCode) = DescribeRefusal(response, settings); + OutputFormatter.WriteError(format, message, suggestion, ExitCodes.CodeFor(exitCode)); + return exitCode; + } + + static (string Message, string Suggestion, int ExitCode) DescribeRefusal(RemoveObserverResponse response, ObserverCommandSettings settings) + { + var @namespace = string.IsNullOrEmpty(response.BlockingNamespace) ? settings.ResolveNamespace() : response.BlockingNamespace; + + return response.Outcome switch + { + ObserverRemovalOutcome.ObserverActive => + ($"Observer '{settings.ObserverId}' is still running in namespace '{@namespace}', so it was not removed.", + "Stop the application that declares this observer, then try again.", + ExitCodes.ValidationError), + + ObserverRemovalOutcome.ObserverSubscribed => + ($"Observer '{settings.ObserverId}' still has a client subscribed to it in namespace '{@namespace}', so it was not removed.", + "Stop the application that declares this observer, then try again.", + ExitCodes.ValidationError), + + ObserverRemovalOutcome.ObserverNotFound => + ($"Observer '{settings.ObserverId}' is not registered in event store '{settings.ResolveEventStore()}', so there was nothing to remove.", + "List the observers to check the identifier: cratis chronicle observers list", + ExitCodes.NotFound), + + _ => + ($"Observer '{settings.ObserverId}' was not removed: {response.Outcome}.", + "Check the observer state with 'cratis chronicle observers show'.", + ExitCodes.ValidationError) + }; + } +} From ebec673c173ddd67ca330f34d4ce1b5a6006db6b Mon Sep 17 00:00:00 2001 From: Einar Date: Fri, 25 Sep 2026 15:13:45 +0200 Subject: [PATCH 4/5] Bump Chronicle to 19.6.1 to pick up observers remove Cratis.Chronicle.Connections, .Contracts and .XUnit.Integration move from 18.1.0 to 19.6.1, the first release carrying RemoveObserver and the RetryPartitionResponse outcome this branch already builds against. 19.5.0 exists as a release tag but was never published - the packages broke consumer startup for an existing external IEventStore implementation, caught and fixed in Chronicle before anything shipped (Cratis/Chronicle#4181). 19.6.1 is the first version after that fix actually reached NuGet. The whole solution builds clean and 1656 specs pass, including the 24 new ones written against the real contract types. --- Directory.Packages.props | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/Directory.Packages.props b/Directory.Packages.props index 230f8de8..416cac5e 100644 --- a/Directory.Packages.props +++ b/Directory.Packages.props @@ -17,8 +17,8 @@ - - + + @@ -68,7 +68,7 @@ - + From 1cacf2e664c777dfad1c2decf99fa98c0875e58d Mon Sep 17 00:00:00 2001 From: Einar Date: Fri, 25 Sep 2026 15:34:26 +0200 Subject: [PATCH 5/5] Update the retry-partition integration spec for the corrected behavior This asserted the exact bug #186 fixed: given a partition key that was never a real failure, it expected "Retry started" and a zero exit code regardless - which is what the command used to say no matter what the kernel answered. Against the live server it now correctly asserts the refusal: a validation exit code, no false claim of a started retry, and the kernel's own message on stderr, where WriteError actually places it. --- .../for_Observers/when_retrying_partition.cs | 16 +++++++++++++--- 1 file changed, 13 insertions(+), 3 deletions(-) diff --git a/Integration/Chronicle/for_Observers/when_retrying_partition.cs b/Integration/Chronicle/for_Observers/when_retrying_partition.cs index 21982f5f..69b457b7 100644 --- a/Integration/Chronicle/for_Observers/when_retrying_partition.cs +++ b/Integration/Chronicle/for_Observers/when_retrying_partition.cs @@ -5,6 +5,13 @@ namespace Cratis.Cli.Integration.Chronicle.for_Observers; +/// +/// Against a real observer with no failed partition of the given key, this is what #186 was about: the command +/// used to await the call, ignore what came back, and print "Retry started" regardless. The key here was never a +/// real failure, so the kernel declines - and the point of the fix is that the command now says so instead of +/// claiming success for a retry that never started. +/// +/// The this specification runs in. [Collection(ChronicleCollection.Name)] public class when_retrying_partition(context context) : CliGiven(context) { @@ -23,9 +30,12 @@ async Task Because() } } - [Fact] void should_return_success_exit_code() => Context.Result.ExitCode.ShouldEqual(ExitCodes.Success); + [Fact] void should_not_return_success_exit_code() => Context.Result.ExitCode.ShouldNotEqual(ExitCodes.Success); - [Fact] void should_contain_retry_started_message() => Context.Result.StandardOutput.ShouldContain("Retry started"); + [Fact] void should_return_a_validation_error_exit_code() => Context.Result.ExitCode.ShouldEqual(ExitCodes.ValidationError); - [Fact] void should_have_no_errors() => Context.Result.StandardError.ShouldEqual(string.Empty); + [Fact] void should_not_claim_a_retry_started() => Context.Result.StandardOutput.ShouldNotContain("Retry started"); + + [Fact] void should_say_the_partition_was_not_among_the_failures() => + Context.Result.StandardError.ShouldContain("is not among the failed partitions"); }