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 @@
-
+
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");
}
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/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.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.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;
- }
-}
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)
+ };
+ }
+}
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'.")
+ };
}