keiro-ops (empty) → 0.12.0.0
raw patch · 21 files changed
+6744/−0 lines, 21 filesdep +aesondep +basedep +bytestring
Dependencies added: aeson, base, bytestring, containers, effectful, effectful-core, hasql, hasql-transaction, hspec, keiro, keiro-migrations, keiro-ops, keiro-pgmq, keiro-test-support, kiroku-store, optparse-applicative, pgmq-migration, process, text, time, uuid, vector
Files
- CHANGELOG.md +85/−0
- LICENSE +28/−0
- app/Main.hs +6/−0
- keiro-ops.cabal +118/−0
- src/Keiro/Ops.hs +217/−0
- src/Keiro/Ops/Embed.hs +32/−0
- src/Keiro/Ops/Env.hs +65/−0
- src/Keiro/Ops/Inbox.hs +235/−0
- src/Keiro/Ops/Outbox.hs +323/−0
- src/Keiro/Ops/Parse.hs +88/−0
- src/Keiro/Ops/Pgmq.hs +233/−0
- src/Keiro/Ops/Projection.hs +183/−0
- src/Keiro/Ops/Rebuild.hs +996/−0
- src/Keiro/Ops/Render.hs +89/−0
- src/Keiro/Ops/ReplayAudit.hs +185/−0
- src/Keiro/Ops/Shard.hs +155/−0
- src/Keiro/Ops/Snapshot.hs +291/−0
- src/Keiro/Ops/Stream.hs +465/−0
- src/Keiro/Ops/Timer.hs +339/−0
- src/Keiro/Ops/Workflow.hs +882/−0
- test/Main.hs +1729/−0
+ CHANGELOG.md view
@@ -0,0 +1,85 @@+# Changelog++All notable changes to `keiro-ops` are recorded here. The format follows+[Keep a Changelog](https://keepachangelog.com/), and the package follows the+[Haskell Package Versioning Policy](https://pvp.haskell.org/).++## Unreleased++## 0.12.0.0 — 2026-08-14++### Breaking Changes++- Requires `kiroku-store >=0.7 && <0.8` for explicit checkpoint lifecycle, the+ public transaction-composable reset API, and the visible-head query used by+ the operator position commands, plus the renewable retention evidence required by+ schema-versioned rebuilds.++### New Features++- Embedded `rebuild reproject-stream GROUP PROJECTION STREAM` previews exact serving+ revision, target, dedup, stream-history, and work-admission facts and executes only+ with `--force`. The positive `--max-events` limit defaults to 1000 and is rechecked+ against locked stream metadata before the group fence. Human and v2 JSON outcomes+ report the admitted limit, cleared rows, replay/application counts, dedup+ inserted/existing counts, and verification; typed refusals have stable operator codes.++- Embedded read-only `rebuild external-read CONTRACT VERSION` inspection and+ preview/`--force` `rebuild retire-external-read CONTRACT VERSION`. Both render the+ supported catalog report with contract state, surface generation, PostgreSQL+ dependents, and execute grants; forced retirement uses the library lifecycle API.+- Embedded `rebuild versioned start|status|resume|abandon`, database-backed+ `rebuild retired`, and preview/`--force` `rebuild drop-retired`. Commands render the+ supported catalog operations reports, derive the physical fleet from the mounted+ validated catalog, and expose stable versioned JSON for runs and retirement blockers.++- `Keiro.Ops.Parse.nonNegativeReader`, a message-parameterized bounded reader+ now shared by global-position, stream-version, and generation options without+ changing their accepted values or domain-specific errors.+- A standalone, schema-checked operations console for Keiro- and Kiroku-owned+ database operations, with human tables and stable JSON generated from the+ same result values.+- Preview-before-`--force` mutations, schema-drift refusal, and typed stream-name+ confirmation for permanent stream operations.+- `AppHooks`, `opsCommandTree`, `runOpsInvocation`, and `mainWithHooks` for+ mounting application-owned workflow resume, timer dispatch, candidate-code+ replay audit, and validated projection-catalog rebuild commands.+- Read-only `stream subscriptions` and+ `projection position --subscription NAME` commands backed by the public+ Kiroku 0.4 durable checkpoint inventory. Both preserve member rows and report+ `global_position_distance`; neither queries Kiroku's private schema or claims+ a relevant-event lag.+- Projection catalog inventory and rebuild JSON expose each subscription's+ stable `checkpointOnMissing` value from the same validated catalog used by+ runtime registration and rebuild planning.+- Embedded catalog operations add `rebuild adopt GROUP...`. Without `--force`+ it classifies every catalog group, shows stored/current slice fingerprints+ and removed groups, and prints the exact force invocation. With `--force` it+ calls the supported transactional adoption API and reports the adopted rows.+ Existing rebuild list and preview tables also expose slice identity.+- Rebuild run tables now include `group_slice`, so status and mutation previews+ expose `$pre-canonical` directly during migration recovery.+- Embedded `wf resume-once` results expose `advanced` and `paced` counts plus+ `sleep_due` and the sorted set of `unregistered_names` in JSON (and the+ corresponding human columns), so an operator can terminate a bounded drain on+ durable progress, identify missing workflow definitions, and distinguish due+ sleeps that require the timer worker rather than another resume pass.++### Bug Fixes++- The non-forced `rebuild adopt` preview now distinguishes the named groups it will adopt+ from out-of-scope catalog drift and warns when skipped groups will still refuse startup+ registration.+- `rebuild status` and the non-forced `rebuild abandon` preview now work for+ pre-canonical runs, enabling the documented abandon, adopt, and fresh-start+ recovery sequence without direct SQL.++### Other Changes++- First public release. Requires `keiro ^>=0.12.0.0`,+ `keiro-migrations ^>=0.12.0.0`, and `keiro-pgmq ^>=0.12.0.0`.+- The source distribution now includes the BSD-3-Clause license file.+- `rebuild adopt` now renders scope-annotated group, registration, and old-name rows and+ reports the forced transaction through `keiro/catalog-adoption-preview/v2` and+ `keiro/catalog-adoption-outcome/v2` JSON envelopes. Preview refuses a requested group+ absent from the catalog with `AdoptGroupNotInCatalog`, matching forced execution.
+ LICENSE view
@@ -0,0 +1,28 @@+BSD 3-Clause License++Copyright (c) 2026, Nadeem Bitar++Redistribution and use in source and binary forms, with or without+modification, are permitted provided that the following conditions are met:++1. Redistributions of source code must retain the above copyright notice, this+ list of conditions and the following disclaimer.++2. Redistributions in binary form must reproduce the above copyright notice,+ this list of conditions and the following disclaimer in the documentation+ and/or other materials provided with the distribution.++3. Neither the name of the copyright holder nor the names of its contributors+ may be used to endorse or promote products derived from this software without+ specific prior written permission.++THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS "AS IS" AND+ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO, THE IMPLIED+WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR PURPOSE ARE+DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT HOLDER OR CONTRIBUTORS BE LIABLE FOR+ANY DIRECT, INDIRECT, INCIDENTAL, SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES+INCLUDING, BUT NOT LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS+OF USE, DATA, OR PROFITS; OR BUSINESS INTERRUPTION HOWEVER CAUSED AND ON ANY+THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT INCLUDING+NEGLIGENCE OR OTHERWISE ARISING IN ANY WAY OUT OF THE USE OF THIS SOFTWARE, EVEN+IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
+ app/Main.hs view
@@ -0,0 +1,6 @@+module Main (main) where++import Keiro.Ops qualified++main :: IO ()+main = Keiro.Ops.main
+ keiro-ops.cabal view
@@ -0,0 +1,118 @@+cabal-version: 3.0+name: keiro-ops+version: 0.12.0.0+synopsis: Operational command-line interface for Keiro deployments+description:+ A standalone database-only console and embeddable command tree for inspecting+ and operating Keiro deployments.++license: BSD-3-Clause+license-file: LICENSE+author: Nadeem Bitar+maintainer: nadeem@gmail.com+copyright: 2026 Nadeem Bitar+category: Operations+homepage: https://github.com/shinzui/keiro#readme+bug-reports: https://github.com/shinzui/keiro/issues+build-type: Simple+tested-with: GHC >=9.12 && <9.13+extra-doc-files: CHANGELOG.md++source-repository head+ type: git+ location: https://github.com/shinzui/keiro.git++common warnings+ ghc-options:+ -Wall -Wcompat -Widentities -Wincomplete-record-updates+ -Wincomplete-uni-patterns -Wpartial-fields -Wredundant-constraints++common shared+ default-language: GHC2024+ default-extensions:+ BlockArguments+ DeriveAnyClass+ DuplicateRecordFields+ ImportQualifiedPost+ LambdaCase+ OverloadedLabels+ OverloadedRecordDot+ OverloadedStrings++library+ import: warnings, shared+ hs-source-dirs: src+ exposed-modules:+ Keiro.Ops+ Keiro.Ops.Embed+ Keiro.Ops.Env+ Keiro.Ops.Inbox+ Keiro.Ops.Outbox+ Keiro.Ops.Parse+ Keiro.Ops.Pgmq+ Keiro.Ops.Projection+ Keiro.Ops.Rebuild+ Keiro.Ops.Render+ Keiro.Ops.ReplayAudit+ Keiro.Ops.Shard+ Keiro.Ops.Snapshot+ Keiro.Ops.Stream+ Keiro.Ops.Timer+ Keiro.Ops.Workflow++ build-depends:+ , aeson >=2.2.2 && <2.3+ , base >=4.21 && <5+ , bytestring >=0.12 && <0.13+ , containers >=0.6 && <0.8+ , effectful >=2.6 && <2.7+ , effectful-core >=2.6 && <2.7+ , hasql >=1.10 && <1.11+ , keiro ^>=0.12.0.0+ , keiro-migrations ^>=0.12.0.0+ , keiro-pgmq ^>=0.12.0.0+ , kiroku-store >=0.7 && <0.8+ , optparse-applicative >=0.18 && <0.20+ , text >=2.1 && <2.2+ , time >=1.12 && <1.15+ , uuid >=1.3 && <1.4+ , vector >=0.13 && <0.14++executable keiro-ops+ import: warnings, shared+ hs-source-dirs: app+ main-is: Main.hs+ ghc-options: -threaded -rtsopts -with-rtsopts=-N+ build-depends:+ , base >=4.21 && <5+ , keiro-ops++test-suite keiro-ops-test+ import: warnings, shared+ type: exitcode-stdio-1.0+ hs-source-dirs: test+ main-is: Main.hs+ ghc-options: -threaded -rtsopts -with-rtsopts=-N+ build-tool-depends: keiro-ops:keiro-ops+ build-depends:+ , aeson >=2.2 && <2.3+ , base >=4.21 && <5+ , bytestring >=0.12 && <0.13+ , containers >=0.6 && <0.8+ , effectful >=2.6 && <2.7+ , effectful-core >=2.6 && <2.7+ , hasql >=1.10 && <1.11+ , hasql-transaction >=1.1 && <1.3+ , hspec >=2.11+ , keiro ^>=0.12.0.0+ , keiro-ops+ , keiro-pgmq ^>=0.12.0.0+ , keiro-test-support+ , kiroku-store >=0.7 && <0.8+ , optparse-applicative >=0.18 && <0.20+ , pgmq-migration >=0.5 && <0.6+ , process >=1.6 && <1.7+ , text >=2.1 && <2.2+ , time >=1.12 && <1.15+ , uuid >=1.3 && <1.4+ , vector >=0.13 && <0.14
+ src/Keiro/Ops.hs view
@@ -0,0 +1,217 @@+module Keiro.Ops+ ( main,+ mainWithHooks,+ AppHooks (..),+ OpsAuditConfig (..),+ emptyAppHooks,+ OpsInvocation,+ opsCommandTree,+ runOpsInvocation,+ )+where++import Control.Exception (SomeException, displayException, fromException, try)+import Data.Foldable (traverse_)+import Data.Maybe (isJust)+import Data.Text qualified as Text+import Data.Text.IO qualified as Text.IO+import Hasql.Connection.Settings qualified as Settings+import Keiro.Migrations.SchemaCheck (renderSchemaDrift, verifyExpectedSchema)+import Keiro.Ops.Embed+import Keiro.Ops.Env+import Keiro.Ops.Inbox qualified as Inbox+import Keiro.Ops.Outbox qualified as Outbox+import Keiro.Ops.Pgmq qualified as Pgmq+import Keiro.Ops.Projection qualified as Projection+import Keiro.Ops.Rebuild qualified as Rebuild+import Keiro.Ops.Render+import Keiro.Ops.ReplayAudit qualified as ReplayAudit+import Keiro.Ops.Shard qualified as Shard+import Keiro.Ops.Snapshot qualified as Snapshot+import Keiro.Ops.Stream qualified as Stream+import Keiro.Ops.Timer qualified as Timer+import Keiro.Ops.Workflow qualified as Workflow+import Kiroku.Store.Connection (defaultConnectionSettings, withStore)+import Options.Applicative+import System.Exit qualified as Exit+import System.IO (stderr)++main :: IO ()+main = mainWithHooks emptyAppHooks++mainWithHooks :: AppHooks -> IO ()+mainWithHooks hooks = do+ invocation <- customExecParser (prefs subparserInline) (opsCommandTree hooks)+ exitCode <- runOpsInvocation hooks invocation+ Exit.exitWith exitCode++runOpsInvocation :: AppHooks -> OpsInvocation -> IO Exit.ExitCode+runOpsInvocation hooks invocation = do+ result <- try (runInvocation hooks invocation)+ case result of+ Left exception ->+ case fromException exception :: Maybe Exit.ExitCode of+ Just exitCode -> pure exitCode+ Nothing -> operationalFailure (Text.pack (displayException (exception :: SomeException)))+ Right exitCode -> pure exitCode++data OpsInvocation = OpsInvocation+ { globalOptions :: !GlobalOptions,+ opsCommand :: !Command+ }++data Command+ = Workflow Workflow.Command+ | Timer Timer.Command+ | Outbox Outbox.Command+ | Inbox Inbox.Command+ | Pgmq Pgmq.Command+ | Projection Projection.Command+ | Shard Shard.Command+ | Snapshot Snapshot.Command+ | Stream Stream.Command+ | ReplayAudit ReplayAudit.Command+ | Rebuild Rebuild.Command++opsCommandTree :: AppHooks -> ParserInfo OpsInvocation+opsCommandTree hooks =+ info+ (invocationParser hooks <**> helper)+ ( fullDesc+ <> progDesc "Inspect and operate a Keiro deployment"+ <> failureCode 2+ )++invocationParser :: AppHooks -> Parser OpsInvocation+invocationParser hooks = OpsInvocation <$> globalOptionsParser <*> commandParser hooks++commandParser :: AppHooks -> Parser Command+commandParser hooks =+ hsubparser+ ( command+ "wf"+ ( info+ (Workflow <$> Workflow.commandParserWithResume (isJust hooks.workflowResume))+ (progDesc "Inspect and operate durable workflows")+ )+ <> command+ "timer"+ ( info+ (Timer <$> Timer.commandParserWithDrain (isJust hooks.timerFire))+ (progDesc "Inspect and operate durable timers")+ )+ <> command+ "outbox"+ (info (Outbox <$> Outbox.commandParser) (progDesc "Inspect and operate the transactional outbox"))+ <> command+ "inbox"+ (info (Inbox <$> Inbox.commandParser) (progDesc "Inspect and operate the integration-event inbox"))+ <> command+ "pgmq"+ (info (Pgmq <$> Pgmq.commandParser) (progDesc "Inspect and operate Keiro PGMQ queues"))+ <> command+ "projection"+ (info (Projection <$> Projection.commandParser) (progDesc "Inspect and operate projection dedup state"))+ <> command+ "shard"+ (info (Shard <$> Shard.commandParser) (progDesc "Inspect and operate sharded-subscription ownership"))+ <> command+ "snapshot"+ (info (Snapshot <$> Snapshot.commandParser) (progDesc "Inspect and operate advisory snapshots"))+ <> command+ "stream"+ (info (Stream <$> Stream.commandParser) (progDesc "Inspect and operate Kiroku streams"))+ <> replayAuditCommand+ <> rebuildCommand+ )+ where+ replayAuditCommand =+ case hooks.replayAudit of+ Nothing -> mempty+ Just _ ->+ command+ "replay-audit"+ (info (ReplayAudit <$> ReplayAudit.commandParser) (progDesc "Audit candidate-code replay against configured targets"))+ rebuildCommand =+ case hooks.projectionCatalog of+ Nothing -> mempty+ Just _ ->+ command+ "rebuild"+ (info (Rebuild <$> Rebuild.commandParser) (progDesc "Inspect and operate the mounted projection catalog"))++runInvocation :: AppHooks -> OpsInvocation -> IO Exit.ExitCode+runInvocation hooks OpsInvocation {globalOptions, opsCommand} = do+ connectionString <- resolveConnectionString globalOptions.databaseUrl+ verified <- verifyExpectedSchema (Settings.connectionString connectionString)+ case verified of+ Left migrationError ->+ operationalFailure ("schema verification failed: " <> Text.pack (show migrationError))+ Right drifts -> do+ let renderedDrifts = map renderSchemaDrift drifts+ traverse_ (Text.IO.hPutStrLn stderr . ("warning: " <>)) renderedDrifts+ if isMutation opsCommand && not (null drifts) && not globalOptions.allowSchemaDrift+ then+ operationalFailure+ "refusing mutation because the live schema differs from this binary; inspect the warnings or pass --allow-schema-drift"+ else withStore (defaultConnectionSettings connectionString) $ \store -> do+ let env =+ OpsEnv+ { store,+ outputMode = globalOptions.outputMode,+ force = globalOptions.force,+ schemaDrift = renderedDrifts,+ allowSchemaDrift = globalOptions.allowSchemaDrift+ }+ runCommand hooks env opsCommand >>= finishOutcome env++isMutation :: Command -> Bool+isMutation = \case+ Workflow workflowCommand -> Workflow.isMutation workflowCommand+ Timer timerCommand -> Timer.isMutation timerCommand+ Outbox outboxCommand -> Outbox.isMutation outboxCommand+ Inbox inboxCommand -> Inbox.isMutation inboxCommand+ Pgmq pgmqCommand -> Pgmq.isMutation pgmqCommand+ Projection projectionCommand -> Projection.isMutation projectionCommand+ Shard shardCommand -> Shard.isMutation shardCommand+ Snapshot snapshotCommand -> Snapshot.isMutation snapshotCommand+ Stream streamCommand -> Stream.isMutation streamCommand+ ReplayAudit _ -> False+ Rebuild rebuildCommand -> Rebuild.isMutation rebuildCommand++runCommand :: AppHooks -> OpsEnv -> Command -> IO OpsOutcome+runCommand hooks env = \case+ Workflow workflowCommand -> Workflow.runCommandWithResume hooks.workflowResume env workflowCommand+ Timer timerCommand -> Timer.runCommandWithFire hooks.timerFire env timerCommand+ Outbox outboxCommand -> Outbox.runCommand env outboxCommand+ Inbox inboxCommand -> Inbox.runCommand env inboxCommand+ Pgmq pgmqCommand -> Pgmq.runCommand env pgmqCommand+ Projection projectionCommand -> Projection.runCommand env projectionCommand+ Shard shardCommand -> Shard.runCommand env shardCommand+ Snapshot snapshotCommand -> Snapshot.runCommand env snapshotCommand+ Stream streamCommand -> Stream.runCommand env streamCommand+ ReplayAudit replayAuditCommand ->+ maybe+ (pure (Failed "replay audit hook is not mounted"))+ (\config -> ReplayAudit.runCommand env config replayAuditCommand)+ hooks.replayAudit+ Rebuild rebuildCommand ->+ maybe+ (pure (Failed "projection catalog hook is not mounted"))+ (\operations -> Rebuild.runCommand env operations rebuildCommand)+ hooks.projectionCatalog++finishOutcome :: OpsEnv -> OpsOutcome -> IO Exit.ExitCode+finishOutcome env = \case+ Succeeded result -> renderResult env result >> pure Exit.ExitSuccess+ SucceededWithExit result exitCode -> renderResult env result >> pure exitCode+ PreviewRequired result reinvocation -> do+ renderResult env result+ Text.IO.hPutStrLn stderr ("preview only; re-run with --force: " <> reinvocation)+ pure (Exit.ExitFailure 1)+ Failed message -> operationalFailure message++operationalFailure :: Text.Text -> IO Exit.ExitCode+operationalFailure message = do+ Text.IO.hPutStrLn stderr ("keiro-ops: " <> message)+ pure (Exit.ExitFailure 1)
+ src/Keiro/Ops/Embed.hs view
@@ -0,0 +1,32 @@+-- | Application-owned hooks for commands that cannot exist in the standalone+-- @keiro-ops@ binary. Commands are mounted only when their hook is present.+--+-- This is the embedding boundary established by+-- @docs/adr/0028-operator-commands-wrap-supported-library-apis-and-respect-schema-ownership.md@.+module Keiro.Ops.Embed+ ( AppHooks (..),+ OpsAuditConfig (..),+ emptyAppHooks,+ )+where++import Keiro.Ops.ReplayAudit (OpsAuditConfig (..))+import Keiro.Ops.Timer (TimerFire)+import Keiro.Ops.Workflow (ResumeHook)+import Keiro.Projection.Catalog.Operations (ProjectionCatalogOperations)++data AppHooks = AppHooks+ { workflowResume :: !(Maybe ResumeHook),+ timerFire :: !(Maybe TimerFire),+ replayAudit :: !(Maybe OpsAuditConfig),+ projectionCatalog :: !(Maybe ProjectionCatalogOperations)+ }++emptyAppHooks :: AppHooks+emptyAppHooks =+ AppHooks+ { workflowResume = Nothing,+ timerFire = Nothing,+ replayAudit = Nothing,+ projectionCatalog = Nothing+ }
+ src/Keiro/Ops/Env.hs view
@@ -0,0 +1,65 @@+module Keiro.Ops.Env+ ( GlobalOptions (..),+ OpsEnv (..),+ OutputMode (..),+ globalOptionsParser,+ resolveConnectionString,+ selectConnectionString,+ )+where++import Data.Text (Text)+import Data.Text qualified as Text+import Kiroku.Store.Connection (KirokuStore)+import Options.Applicative+import System.Environment (lookupEnv)++data OutputMode+ = HumanTable+ | Json+ deriving stock (Eq, Show)++data GlobalOptions = GlobalOptions+ { databaseUrl :: !(Maybe Text),+ outputMode :: !OutputMode,+ force :: !Bool,+ allowSchemaDrift :: !Bool+ }+ deriving stock (Eq, Show)++data OpsEnv = OpsEnv+ { store :: !KirokuStore,+ outputMode :: !OutputMode,+ force :: !Bool,+ schemaDrift :: ![Text],+ allowSchemaDrift :: !Bool+ }++globalOptionsParser :: Parser GlobalOptions+globalOptionsParser =+ GlobalOptions+ <$> optional+ ( Text.pack+ <$> strOption+ ( long "database-url"+ <> metavar "URL"+ <> help+ "PostgreSQL URI or keyword/value string; defaults to KEIRO_OPS_DATABASE_URL, DATABASE_URL, then libpq PG* variables"+ )+ )+ <*> flag HumanTable Json (long "json" <> help "Emit machine-readable JSON")+ <*> switch (long "force" <> help "Apply a mutating command after its preview")+ <*> switch+ ( long "allow-schema-drift"+ <> help "Allow a mutating command despite a failed schema agreement check"+ )++resolveConnectionString :: Maybe Text -> IO Text+resolveConnectionString explicit = do+ keiroOpsDatabaseUrl <- fmap Text.pack <$> lookupEnv "KEIRO_OPS_DATABASE_URL"+ databaseUrl <- fmap Text.pack <$> lookupEnv "DATABASE_URL"+ pure (selectConnectionString explicit keiroOpsDatabaseUrl databaseUrl)++selectConnectionString :: Maybe Text -> Maybe Text -> Maybe Text -> Text+selectConnectionString explicit keiroOpsDatabaseUrl databaseUrl =+ maybe "" id (explicit <|> keiroOpsDatabaseUrl <|> databaseUrl)
+ src/Keiro/Ops/Inbox.hs view
@@ -0,0 +1,235 @@+-- | Operational adapters for the Keiro inbox.+--+-- All state changes flow through 'Keiro.Inbox' as required by ADR 28.+module Keiro.Ops.Inbox+ ( Command (..),+ ListOptions (..),+ commandParser,+ isMutation,+ runCommand,+ )+where++import Data.Aeson (Value, object, (.=))+import Data.Aeson qualified as Aeson+import Data.Int (Int64)+import Data.Text (Text)+import Data.Text qualified as Text+import Data.Text.Encoding qualified as Text.Encoding+import Data.Text.Encoding.Error qualified as Text.Error+import Data.Time (NominalDiffTime, UTCTime, getCurrentTime)+import Data.UUID qualified as UUID+import Effectful (Eff, IOE)+import Effectful.Error.Static (Error)+import Keiro.Inbox+import Keiro.Integration.Event (IntegrationEvent (..))+import Keiro.Ops.Env (OpsEnv (..), OutputMode (..))+import Keiro.Ops.Parse (durationReader, positiveIntReader)+import Keiro.Ops.Render+import Kiroku.Store.Effect (Store, runStoreIO)+import Kiroku.Store.Error (StoreError)+import Kiroku.Store.Transaction (runTransaction)+import Kiroku.Store.Types (EventId (..), GlobalPosition (..))+import Options.Applicative hiding (action, value)+import Options.Applicative qualified as Opt++data ListOptions = ListOptions+ { source :: !Text,+ status :: !(Maybe InboxStatus),+ limit :: !Int+ }+ deriving stock (Eq, Show)++data Command+ = Backlog+ | List !ListOptions+ | Show !Text !Text+ | Gc !NominalDiffTime+ | MarkFailed !Text !Text !Text+ deriving stock (Eq, Show)++commandParser :: Parser Command+commandParser =+ hsubparser+ ( command "backlog" (info (pure Backlog) (progDesc "Count processing and failed inbox rows"))+ <> command "list" (info listParser (progDesc "List inbox rows for a source"))+ <> command "show" (info showParser (progDesc "Inspect one inbox row"))+ <> command "gc" (info gcParser (progDesc "Preview or delete retained completed rows"))+ <> command "mark-failed" (info markFailedParser (progDesc "Preview or mark an inbox row permanently failed"))+ )+ where+ listParser =+ List+ <$> ( ListOptions+ <$> textOption "source" "SOURCE" "Producing bounded-context source"+ <*> optional (option statusReader (long "status" <> metavar "STATUS" <> help "processing, completed, or failed"))+ <*> option positiveIntReader (long "limit" <> metavar "N" <> Opt.value 100 <> showDefault <> help "Maximum rows")+ )+ showParser =+ Show+ <$> argument (Text.pack <$> str) (metavar "SOURCE")+ <*> argument (Text.pack <$> str) (metavar "MESSAGE_ID")+ gcParser =+ Gc+ <$> option durationReader (long "older-than" <> metavar "DURATION" <> Opt.value 2592000 <> showDefaultWith (const "30d") <> help "Completed-row retention age")+ markFailedParser =+ MarkFailed+ <$> argument (Text.pack <$> str) (metavar "SOURCE")+ <*> argument (Text.pack <$> str) (metavar "MESSAGE_ID")+ <*> textOption "reason" "TEXT" "Permanent-failure reason"++textOption :: String -> String -> String -> Parser Text+textOption name metavarText helpText = Text.pack <$> strOption (long name <> metavar metavarText <> help helpText)++statusReader :: ReadM InboxStatus+statusReader = eitherReader (either (Left . Text.unpack) Right . parseInboxStatus . Text.pack)++isMutation :: Command -> Bool+isMutation = \case+ Backlog -> False+ List {} -> False+ Show {} -> False+ Gc {} -> True+ MarkFailed {} -> True++runCommand :: OpsEnv -> Command -> IO OpsOutcome+runCommand env = \case+ Backlog -> runAction env countInboxBacklog (Succeeded . countResult "inbox_backlog")+ List options ->+ runAction env (listInbox options.source) $ \rows ->+ Succeeded (inboxListResult (take options.limit (filterByStatus options.status rows)))+ Show source messageId ->+ runAction env (lookupInbox source messageId) (Succeeded . maybe emptyResult (inboxListResult . pure))+ Gc olderThan -> runGc env olderThan+ MarkFailed source messageId reason -> runMarkFailed env source messageId reason++filterByStatus :: Maybe InboxStatus -> [InboxRow] -> [InboxRow]+filterByStatus Nothing = id+filterByStatus (Just expected) = filter ((== expected) . (.status))++runGc :: OpsEnv -> NominalDiffTime -> IO OpsOutcome+runGc env olderThan = do+ now <- getCurrentTime+ if env.force+ then runAction env (garbageCollectCompleted olderThan now) (Succeeded . countResult "deleted")+ else runAction env (listCompletedInboxGcCandidates olderThan now) $ \rows ->+ PreviewRequired+ (inboxListResult rows)+ (forceInvocation env ["inbox", "gc", "--older-than", durationText olderThan])++runMarkFailed :: OpsEnv -> Text -> Text -> Text -> IO OpsOutcome+runMarkFailed env source messageId reason+ | not env.force =+ runAction env (lookupInbox source messageId) $ \row ->+ PreviewRequired+ (markFailedPreview source messageId row)+ (forceInvocation env ["inbox", "mark-failed", source, messageId, "--reason", reason])+ | otherwise = do+ now <- getCurrentTime+ runAction+ env+ ( do+ runTransaction (markFailedTx source messageId reason now)+ lookupInbox source messageId+ )+ $ \row ->+ Succeeded (maybe emptyResult (inboxListResult . pure) row)++runAction :: OpsEnv -> Eff '[Store, Error StoreError, IOE] a -> (a -> OpsOutcome) -> IO OpsOutcome+runAction env action onSuccess = do+ result <- runStoreIO env.store action+ pure $ either (Failed . Text.pack . show) onSuccess result++countResult :: Text -> Int -> OpsResult+countResult label count = OpsResult [label] [[showText count]] (object ["metric" .= label, "count" .= count])++inboxListResult :: [InboxRow] -> OpsResult+inboxListResult inboxRows =+ OpsResult+ { headers = ["source", "message_id", "status", "attempts", "received_at", "last_error"],+ rows = map inboxRow inboxRows,+ jsonValue = Aeson.toJSON (map inboxJson inboxRows)+ }++inboxRow :: InboxRow -> [Text]+inboxRow row =+ [ row.source,+ row.dedupeKey,+ inboxStatusText row.status,+ showText row.attemptCount,+ timeText row.receivedAt,+ maybe "" (truncateCell 120) row.lastError+ ]++inboxJson :: InboxRow -> Value+inboxJson row =+ object+ [ "source" .= row.source,+ "dedupe_key" .= row.dedupeKey,+ "event" .= eventJson row.event,+ "kafka" .= fmap kafkaJson row.kafka,+ "status" .= inboxStatusText row.status,+ "attempt_count" .= row.attemptCount,+ "received_at" .= row.receivedAt,+ "completed_at" .= row.completedAt,+ "failed_at" .= row.failedAt,+ "last_error" .= row.lastError+ ]++eventJson :: IntegrationEvent -> Value+eventJson event =+ object+ [ "message_id" .= event.messageId,+ "source" .= event.source,+ "destination" .= event.destination,+ "key" .= event.key,+ "event_type" .= event.eventType,+ "schema_version" .= event.schemaVersion,+ "source_event_id" .= fmap eventIdText event.sourceEventId,+ "source_global_position" .= fmap globalPositionInt event.sourceGlobalPosition,+ "payload" .= payloadValue event,+ "occurred_at" .= event.occurredAt+ ]++kafkaJson :: KafkaDeliveryRef -> Value+kafkaJson ref = object ["topic" .= ref.topic, "partition" .= ref.partition, "offset" .= ref.offset]++payloadValue :: IntegrationEvent -> Value+payloadValue event =+ either+ (const (Aeson.String (Text.Encoding.decodeUtf8With Text.Error.lenientDecode event.payloadBytes)))+ id+ (Aeson.eitherDecodeStrict' event.payloadBytes)++markFailedPreview :: Text -> Text -> Maybe InboxRow -> OpsResult+markFailedPreview source messageId row =+ OpsResult+ { headers = ["source", "message_id", "current_status", "disposition"],+ rows = [[source, messageId, maybe "not_found" (inboxStatusText . (.status)) row, disposition]],+ jsonValue = object ["preview" .= True, "disposition" .= disposition, "inbox" .= fmap inboxJson row]+ }+ where+ disposition = maybe "not_found" (const "would_mark_failed") row++eventIdText :: EventId -> Text+eventIdText (EventId value) = UUID.toText value++globalPositionInt :: GlobalPosition -> Int64+globalPositionInt (GlobalPosition value) = value++timeText :: UTCTime -> Text+timeText = Text.pack . show++showText :: (Show a) => a -> Text+showText = Text.pack . show++durationText :: NominalDiffTime -> Text+durationText = showText . (realToFrac :: NominalDiffTime -> Double)++forceInvocation :: OpsEnv -> [Text] -> Text+forceInvocation env arguments = Text.unwords (map shellQuote ("keiro-ops" : arguments <> globalFlags <> ["--force"]))+ where+ globalFlags = ["--json" | env.outputMode == Json] <> ["--allow-schema-drift" | env.allowSchemaDrift]++shellQuote :: Text -> Text+shellQuote value = "'" <> Text.replace "'" "'\"'\"'" value <> "'"
+ src/Keiro/Ops/Outbox.hs view
@@ -0,0 +1,323 @@+-- | Operational adapters for the Keiro outbox and dispatch dead letters.+--+-- Mutations and previews use only the public owning-library operations required+-- by ADR 28; this module never reaches into either schema directly.+module Keiro.Ops.Outbox+ ( Command (..),+ ListOptions (..),+ commandParser,+ isMutation,+ runCommand,+ )+where++import Data.Aeson (Value, object, (.=))+import Data.Aeson qualified as Aeson+import Data.Int (Int64)+import Data.Text (Text)+import Data.Text qualified as Text+import Data.Text.Encoding qualified as Text.Encoding+import Data.Text.Encoding.Error qualified as Text.Error+import Data.Time (NominalDiffTime, UTCTime, getCurrentTime)+import Data.UUID qualified as UUID+import Effectful (Eff, IOE)+import Effectful.Error.Static (Error)+import Keiro.DeadLetter+import Keiro.Integration.Event (IntegrationEvent (..))+import Keiro.Ops.Env (OpsEnv (..), OutputMode (..))+import Keiro.Ops.Parse (durationReader, positiveIntReader)+import Keiro.Ops.Render+import Keiro.Outbox+import Kiroku.Store.Effect (Store, runStoreIO)+import Kiroku.Store.Error (StoreError)+import Kiroku.Store.Types (EventId (..), GlobalPosition (..), StreamName (..))+import Options.Applicative hiding (action, value)+import Options.Applicative qualified as Opt++data ListOptions = ListOptions+ { source :: !Text,+ status :: !(Maybe OutboxStatus),+ destination :: !(Maybe Text),+ limit :: !Int+ }+ deriving stock (Eq, Show)++data Command+ = Backlog+ | List !ListOptions+ | Show !OutboxId+ | RequeueStuck !NominalDiffTime !Int+ | GcSent !NominalDiffTime+ | MaintenancePass+ | DispatchDeadLetters !Text !Int+ deriving stock (Eq, Show)++commandParser :: Parser Command+commandParser =+ hsubparser+ ( command "backlog" (info (pure Backlog) (progDesc "Count claimable outbox rows"))+ <> command "list" (info listParser (progDesc "List outbox rows for a source"))+ <> command "show" (info (Show <$> outboxIdArgument) (progDesc "Inspect one outbox row"))+ <> command "requeue-stuck" (info requeueParser (progDesc "Preview or reclaim stale publishing rows"))+ <> command "gc-sent" (info gcParser (progDesc "Preview or delete retained sent rows"))+ <> command "maintenance-pass" (info (pure MaintenancePass) (progDesc "Preview or run one default outbox maintenance pass"))+ <> command+ "dead-letters"+ (info (hsubparser (command "list" (info deadLettersParser (progDesc "List rejected process-manager or router dispatches")))) (progDesc "Inspect rejected process-manager or router dispatches"))+ )+ where+ listParser =+ List+ <$> ( ListOptions+ <$> textOption "source" "SOURCE" "Producing bounded-context source"+ <*> optional (option statusReader (long "status" <> metavar "STATUS" <> help "pending, publishing, sent, failed, or dead"))+ <*> optional (textOption "destination" "DESTINATION" "Destination filter")+ <*> option positiveIntReader (long "limit" <> metavar "N" <> Opt.value 100 <> showDefault <> help "Maximum rows")+ )+ requeueParser =+ RequeueStuck+ <$> option durationReader (long "older-than" <> metavar "DURATION" <> Opt.value 300 <> showDefaultWith (const "5m") <> help "Minimum publishing age")+ <*> option positiveIntReader (long "max-attempts" <> metavar "N" <> Opt.value 10 <> showDefault <> help "Attempt ceiling; exhausted rows become dead")+ gcParser =+ GcSent+ <$> option durationReader (long "older-than" <> metavar "DURATION" <> Opt.value 2592000 <> showDefaultWith (const "30d") <> help "Sent-row retention age")+ deadLettersParser =+ DispatchDeadLetters+ <$> textOption "dispatcher" "NAME" "Process-manager or router dispatcher name"+ <*> option positiveIntReader (long "limit" <> metavar "N" <> Opt.value 100 <> showDefault <> help "Maximum rows")++textOption :: String -> String -> String -> Parser Text+textOption name metavarText helpText =+ Text.pack <$> strOption (long name <> metavar metavarText <> help helpText)++outboxIdArgument :: Parser OutboxId+outboxIdArgument =+ OutboxId <$> argument uuidReader (metavar "OUTBOX_ID")++uuidReader :: ReadM UUID.UUID+uuidReader = eitherReader $ \raw -> maybe (Left "expected a UUID") Right (UUID.fromString raw)++statusReader :: ReadM OutboxStatus+statusReader = eitherReader (firstText . parseStatus . Text.pack)+ where+ firstText = either (Left . Text.unpack) Right++isMutation :: Command -> Bool+isMutation = \case+ Backlog -> False+ List {} -> False+ Show {} -> False+ DispatchDeadLetters {} -> False+ RequeueStuck {} -> True+ GcSent {} -> True+ MaintenancePass -> True++runCommand :: OpsEnv -> Command -> IO OpsOutcome+runCommand env = \case+ Backlog -> runAction env countOutboxBacklog (Succeeded . countResult "outbox_backlog")+ List options -> runAction env (listOutbox options.source) (Succeeded . outboxListResult . applyListOptions options)+ Show outboxId -> runAction env (lookupOutbox outboxId) (Succeeded . maybe emptyResult (outboxListResult . pure))+ RequeueStuck olderThan maxAttempts -> runRequeue env olderThan maxAttempts+ GcSent olderThan -> runGc env olderThan+ MaintenancePass -> runMaintenance env+ DispatchDeadLetters dispatcher limit ->+ runAction env (listDispatchDeadLetters dispatcher) (Succeeded . dispatchListResult . take limit)++applyListOptions :: ListOptions -> [OutboxRow] -> [OutboxRow]+applyListOptions options =+ take options.limit+ . filter (maybe (const True) (\expected row -> row.status == expected) options.status)+ . filter (maybe (const True) (\expected row -> row.event.destination == expected) options.destination)++runRequeue :: OpsEnv -> NominalDiffTime -> Int -> IO OpsOutcome+runRequeue env olderThan maxAttempts = do+ now <- getCurrentTime+ if env.force+ then runAction env (requeueStuckOutbox maxAttempts olderThan now) $ \(requeued, deadLettered) ->+ Succeeded+ OpsResult+ { headers = ["requeued", "dead_lettered"],+ rows = [[showText requeued, showText deadLettered]],+ jsonValue = object ["requeued" .= requeued, "dead_lettered" .= deadLettered]+ }+ else runAction env (listStuckOutbox olderThan now) $ \rows ->+ PreviewRequired+ (outboxPreviewResult maxAttempts rows)+ (forceInvocation env ["outbox", "requeue-stuck", "--older-than", durationText olderThan, "--max-attempts", showText maxAttempts])++runGc :: OpsEnv -> NominalDiffTime -> IO OpsOutcome+runGc env olderThan = do+ now <- getCurrentTime+ if env.force+ then runAction env (garbageCollectSent olderThan now) (Succeeded . countResult "deleted")+ else runAction env (listSentOutboxGcCandidates olderThan now) $ \rows ->+ PreviewRequired+ (outboxListResult rows)+ (forceInvocation env ["outbox", "gc-sent", "--older-than", durationText olderThan])++runMaintenance :: OpsEnv -> IO OpsOutcome+runMaintenance env = do+ now <- getCurrentTime+ let options = defaultMaintenanceOptions+ if env.force+ then runAction env (outboxMaintenancePass options Nothing) $ \summary ->+ Succeeded+ OpsResult+ { headers = ["requeued", "dead_lettered", "backlog"],+ rows = [[showText summary.requeued, showText summary.deadLettered, showText summary.backlog]],+ jsonValue = object ["requeued" .= summary.requeued, "dead_lettered" .= summary.deadLettered, "backlog" .= summary.backlog]+ }+ else runAction env (listStuckOutbox options.publishingTimeout now) $ \rows ->+ PreviewRequired+ (outboxPreviewResult options.maxAttempts rows)+ (forceInvocation env ["outbox", "maintenance-pass"])++runAction :: OpsEnv -> Eff '[Store, Error StoreError, IOE] a -> (a -> OpsOutcome) -> IO OpsOutcome+runAction env action onSuccess = do+ result <- runStoreIO env.store action+ pure $ either (Failed . Text.pack . show) onSuccess result++countResult :: Text -> Int -> OpsResult+countResult label count =+ OpsResult [label] [[showText count]] (object ["count" .= count, "metric" .= label])++outboxListResult :: [OutboxRow] -> OpsResult+outboxListResult outboxRows =+ OpsResult+ { headers = ["id", "source", "destination", "status", "attempts", "created_at", "last_error"],+ rows = map outboxRow outboxRows,+ jsonValue = Aeson.toJSON (map outboxJson outboxRows)+ }++outboxRow :: OutboxRow -> [Text]+outboxRow row =+ [ outboxIdText row.outboxId,+ row.event.source,+ row.event.destination,+ statusText row.status,+ showText row.attemptCount,+ timeText row.createdAt,+ maybe "" (truncateCell 120) row.lastError+ ]++outboxJson :: OutboxRow -> Value+outboxJson row =+ object+ [ "outbox_id" .= outboxIdText row.outboxId,+ "message_id" .= row.event.messageId,+ "source" .= row.event.source,+ "destination" .= row.event.destination,+ "key" .= row.event.key,+ "event_type" .= row.event.eventType,+ "schema_version" .= row.event.schemaVersion,+ "source_event_id" .= fmap eventIdText row.event.sourceEventId,+ "source_global_position" .= fmap globalPositionInt row.event.sourceGlobalPosition,+ "payload" .= payloadValue row.event,+ "status" .= statusText row.status,+ "attempt_count" .= row.attemptCount,+ "next_attempt_at" .= row.nextAttemptAt,+ "last_error" .= row.lastError,+ "published_at" .= row.publishedAt,+ "created_at" .= row.createdAt,+ "updated_at" .= row.updatedAt+ ]++payloadValue :: IntegrationEvent -> Value+payloadValue event =+ either+ (const (Aeson.String (Text.Encoding.decodeUtf8With Text.Error.lenientDecode event.payloadBytes)))+ id+ (Aeson.eitherDecodeStrict' event.payloadBytes)++outboxPreviewResult :: Int -> [OutboxRow] -> OpsResult+outboxPreviewResult maxAttempts outboxRows =+ OpsResult+ { headers = ["id", "current_status", "disposition", "attempts"],+ rows =+ [ [outboxIdText row.outboxId, statusText row.status, disposition row, showText row.attemptCount]+ | row <- outboxRows+ ],+ jsonValue =+ Aeson.toJSON+ [ object ["outbox" .= outboxJson row, "disposition" .= disposition row]+ | row <- outboxRows+ ]+ }+ where+ disposition row+ | row.attemptCount >= maxAttempts = "would_dead_letter"+ | otherwise = "would_requeue"++dispatchListResult :: [DispatchDeadLetterRecord] -> OpsResult+dispatchListResult records =+ OpsResult+ { headers = ["id", "kind", "dispatcher", "correlation", "target", "error", "attempts", "created_at"],+ rows = map dispatchRow records,+ jsonValue = Aeson.toJSON (map dispatchJson records)+ }++dispatchRow :: DispatchDeadLetterRecord -> [Text]+dispatchRow row =+ [ showText row.deadLetterId,+ dispatcherKindText row.dispatcherKind,+ row.dispatcherName,+ row.correlationId,+ streamNameText row.targetStreamName,+ row.errorClass <> ": " <> truncateCell 100 row.errorDetail,+ showText row.attemptCount,+ timeText row.createdAt+ ]++dispatchJson :: DispatchDeadLetterRecord -> Value+dispatchJson row =+ object+ [ "dead_letter_id" .= row.deadLetterId,+ "dispatcher_kind" .= dispatcherKindText row.dispatcherKind,+ "dispatcher_name" .= row.dispatcherName,+ "correlation_id" .= row.correlationId,+ "source_event_id" .= eventIdText row.sourceEventId,+ "source_global_position" .= globalPositionInt row.sourceGlobalPosition,+ "emit_index" .= row.emitIndex,+ "target_stream_name" .= streamNameText row.targetStreamName,+ "error_class" .= row.errorClass,+ "error_detail" .= row.errorDetail,+ "attempt_count" .= row.attemptCount,+ "created_at" .= row.createdAt+ ]++dispatcherKindText :: DispatcherKind -> Text+dispatcherKindText = \case+ DispatcherProcessManager -> "process_manager"+ DispatcherRouter -> "router"++outboxIdText :: OutboxId -> Text+outboxIdText (OutboxId value) = UUID.toText value++eventIdText :: EventId -> Text+eventIdText (EventId value) = UUID.toText value++globalPositionInt :: GlobalPosition -> Int64+globalPositionInt (GlobalPosition value) = value++streamNameText :: StreamName -> Text+streamNameText (StreamName value) = value++timeText :: UTCTime -> Text+timeText = Text.pack . show++showText :: (Show a) => a -> Text+showText = Text.pack . show++durationText :: NominalDiffTime -> Text+durationText = showText . (realToFrac :: NominalDiffTime -> Double)++forceInvocation :: OpsEnv -> [Text] -> Text+forceInvocation env arguments =+ Text.unwords (map shellQuote ("keiro-ops" : arguments <> globalFlags <> ["--force"]))+ where+ globalFlags =+ ["--json" | env.outputMode == Json]+ <> ["--allow-schema-drift" | env.allowSchemaDrift]++shellQuote :: Text -> Text+shellQuote value = "'" <> Text.replace "'" "'\"'\"'" value <> "'"
+ src/Keiro/Ops/Parse.hs view
@@ -0,0 +1,88 @@+module Keiro.Ops.Parse+ ( durationReader,+ parseDuration,+ readBoundedIntegral,+ positiveIntReader,+ nonNegativeReader,+ nonNegativeIntReader,+ )+where++import Data.Char (toLower)+import Data.Time (NominalDiffTime)+import Options.Applicative (ReadM, eitherReader)+import Text.Read qualified as Read++durationReader :: ReadM NominalDiffTime+durationReader = eitherReader parseDuration++parseDuration :: String -> Either String NominalDiffTime+parseDuration input = do+ let (numberText, multiplier) =+ case reverse input of+ suffix : rest+ | Just factor <- durationFactor (toLower suffix) ->+ (reverse rest, factor)+ _ -> (input, 1)+ value <- maybe (Left malformed) Right (Read.readMaybe numberText :: Maybe Double)+ let scaled = value * multiplier+ if isNaN scaled || isInfinite scaled || scaled < 0+ then Left malformed+ else+ if scaled > maxDurationSeconds+ then Left tooLarge+ else Right (realToFrac scaled)+ where+ malformed =+ "invalid duration "+ <> show input+ <> ": expected a finite, non-negative number of seconds, optionally with an s, m, h, or d suffix"+ tooLarge =+ "invalid duration "+ <> show input+ <> ": exceeds the maximum supported duration of 9.0e12 seconds (about 285000 years)"++-- | Upper bound on any operator-supplied duration, in seconds. PostgreSQL's+-- binary timestamptz format is Int64 microseconds since 2000-01-01 (maximum+-- about 9.22e12 seconds); a larger duration wraps modulo 2^64 into an arbitrary+-- cutoff. 9.0e12 seconds is comfortably inside that range and far beyond any+-- legitimate retention.+maxDurationSeconds :: Double+maxDurationSeconds = 9.0e12++-- | Parse through unbounded 'Integer' and admit the value only when it fits the+-- requested bounded integral type. Reading directly at a bounded type silently+-- wraps oversized literals.+readBoundedIntegral :: forall a. (Integral a, Bounded a) => String -> Maybe a+readBoundedIntegral raw =+ case reads raw :: [(Integer, String)] of+ [(value, "")]+ | value >= toInteger (minBound :: a),+ value <= toInteger (maxBound :: a) ->+ Just (fromInteger value)+ _ -> Nothing++positiveIntReader :: ReadM Int+positiveIntReader = eitherReader $ \raw ->+ case readBoundedIntegral raw of+ Just n | n > 0 -> Right n+ _ -> Left "expected a positive integer"++nonNegativeIntReader :: ReadM Int+nonNegativeIntReader = nonNegativeReader "expected a non-negative integer"++-- | A bounded, non-negative integral reader whose failure message names the+-- domain concept being parsed (global position, stream version, generation).+nonNegativeReader :: forall a. (Integral a, Bounded a) => String -> ReadM a+nonNegativeReader message = eitherReader $ \raw ->+ case readBoundedIntegral raw of+ Just value | value >= 0 -> Right value+ _ -> Left message++durationFactor :: Char -> Maybe Double+durationFactor = \case+ 's' -> Just 1+ 'm' -> Just 60+ 'h' -> Just 3600+ 'd' -> Just 86400+ _ -> Nothing
+ src/Keiro/Ops/Pgmq.hs view
@@ -0,0 +1,233 @@+-- | Operational adapters for Keiro PGMQ dead-letter queues.+--+-- The module uses the versioned Keiro PGMQ DLQ helpers and never queries PGMQ+-- tables itself, preserving ADR 1's envelope contract and ADR 28's ownership+-- boundary.+module Keiro.Ops.Pgmq+ ( Command (..),+ DlqCommand (..),+ commandParser,+ isMutation,+ runCommand,+ )+where++import Data.Aeson (Value, object, (.=))+import Data.Aeson qualified as Aeson+import Data.Int (Int64)+import Data.Text (Text)+import Data.Text qualified as Text+import Data.Text.IO qualified as Text.IO+import Keiro.Ops.Env (OpsEnv (..), OutputMode (..))+import Keiro.Ops.Parse (positiveIntReader, readBoundedIntegral)+import Keiro.Ops.Render+import Keiro.PGMQ+import Kiroku.Store.Connection (KirokuStore (..))+import Options.Applicative hiding (action, value)+import Options.Applicative qualified as Opt+import System.IO (hFlush, stdout)++data Command = Dlq !DlqCommand+ deriving stock (Eq, Show)++data DlqCommand+ = Read !Text !Int+ | Redrive !Text !Int+ | Archive !Text !(Maybe Int64) !Int+ | Purge !Text+ deriving stock (Eq, Show)++commandParser :: Parser Command+commandParser =+ Dlq+ <$> hsubparser+ (command "dlq" (info dlqParser (progDesc "Inspect and operate a Keiro PGMQ dead-letter queue")))++dlqParser :: Parser DlqCommand+dlqParser =+ hsubparser+ ( command "read" (info readParser (progDesc "Read and decode visible DLQ entries"))+ <> command "redrive" (info redriveParser (progDesc "Preview or move DLQ entries back to the main queue"))+ <> command "archive" (info archiveParser (progDesc "Preview or archive DLQ entries for retention"))+ <> command "purge" (info purgeParser (progDesc "Preview or permanently purge a DLQ"))+ )+ where+ readParser = Read <$> queueOption <*> limitOption 20+ redriveParser = Redrive <$> queueOption <*> limitOption 100+ archiveParser =+ Archive+ <$> queueOption+ <*> optional (option int64Reader (long "entry" <> metavar "MESSAGE_ID" <> help "Archive only this DLQ message id"))+ <*> limitOption 100+ purgeParser = Purge <$> queueOption++queueOption :: Parser Text+queueOption = Text.pack <$> strOption (long "queue" <> metavar "QUEUE" <> help "Logical Keiro job queue name")++limitOption :: Int -> Parser Int+limitOption defaultLimit = option positiveIntReader (long "limit" <> metavar "N" <> Opt.value defaultLimit <> showDefault <> help "Maximum entries")++int64Reader :: ReadM Int64+int64Reader = eitherReader $ \raw ->+ case readBoundedIntegral raw of+ Just n | n > 0 -> Right n+ _ -> Left "expected a positive message id"++isMutation :: Command -> Bool+isMutation (Dlq dlqCommand) = case dlqCommand of+ Read {} -> False+ Redrive {} -> True+ Archive {} -> True+ Purge {} -> True++runCommand :: OpsEnv -> Command -> IO OpsOutcome+runCommand env (Dlq dlqCommand) = case dlqCommand of+ Read queue limit -> handlePgmq (runJobEff (pgmqRuntime env) (readDlq (rawJob queue) (fromIntegral limit))) (Succeeded . dlqListResult queue)+ Redrive queue limit -> runRedrive env queue limit+ Archive queue entry limit -> runArchive env queue entry limit+ Purge queue -> runPurge env queue++runRedrive :: OpsEnv -> Text -> Int -> IO OpsOutcome+runRedrive env queue limit+ | env.force =+ handlePgmq (runJobEff (pgmqRuntime env) (redriveDlq (rawJob queue) limit)) $ \moved ->+ Succeeded (mutationCountResult "redrive" queue moved)+ | otherwise = previewFromDepth env "redrive" queue (Just limit) ["pgmq", "dlq", "redrive", "--queue", queue, "--limit", showText limit]++runArchive :: OpsEnv -> Text -> Maybe Int64 -> Int -> IO OpsOutcome+runArchive env queue entry limit+ | env.force = case entry of+ Just messageId ->+ handlePgmq (runJobEff (pgmqRuntime env) (archiveDlqEntryById (rawJob queue) messageId)) $ \archived ->+ Succeeded+ OpsResult+ { headers = ["operation", "queue", "message_id", "archived"],+ rows = [["archive", queue, showText messageId, boolText archived]],+ jsonValue = object ["operation" .= ("archive" :: Text), "queue" .= queue, "message_id" .= messageId, "archived" .= archived]+ }+ Nothing ->+ handlePgmq (runJobEff (pgmqRuntime env) (archiveDlq (rawJob queue) limit)) $ \archived ->+ Succeeded (mutationCountResult "archive" queue archived)+ | otherwise =+ case entry of+ Just messageId ->+ pure+ ( PreviewRequired+ OpsResult+ { headers = ["operation", "queue", "message_id", "disposition"],+ rows = [["archive", queue, showText messageId, "would_archive_if_present"]],+ jsonValue = object ["preview" .= True, "operation" .= ("archive" :: Text), "queue" .= queue, "message_id" .= messageId]+ }+ (forceInvocation env ["pgmq", "dlq", "archive", "--queue", queue, "--entry", showText messageId])+ )+ Nothing -> previewFromDepth env "archive" queue (Just limit) ["pgmq", "dlq", "archive", "--queue", queue, "--limit", showText limit]++runPurge :: OpsEnv -> Text -> IO OpsOutcome+runPurge env queue+ | env.force = do+ confirmed <- confirmPurge env queue+ if confirmed+ then handlePgmq (runJobEff (pgmqRuntime env) (purgeDlq (rawJob queue))) $ \() ->+ Succeeded (messageResult ("purged DLQ for " <> queue))+ else pure (Failed "queue-name confirmation did not match; DLQ purge cancelled")+ | otherwise = previewFromDepth env "purge" queue Nothing ["pgmq", "dlq", "purge", "--queue", queue]++confirmPurge :: OpsEnv -> Text -> IO Bool+confirmPurge env queue+ | env.outputMode == Json = pure True+ | otherwise = do+ Text.IO.putStr ("type the queue name to confirm: " <> queue <> "\n> ")+ hFlush stdout+ entered <- Text.IO.getLine+ pure (entered == queue)++previewFromDepth :: OpsEnv -> Text -> Text -> Maybe Int -> [Text] -> IO OpsOutcome+previewFromDepth env operation queue requested arguments =+ handlePgmq (runJobEff (pgmqRuntime env) (jobDlqMetrics (rawJob queue))) $ \metrics ->+ let affected = maybe metrics.queueLength (min metrics.queueLength . fromIntegral) requested+ in PreviewRequired+ OpsResult+ { headers = ["operation", "queue", "available", "would_affect"],+ rows = [[operation, queue, showText metrics.queueLength, showText affected]],+ jsonValue =+ object+ [ "preview" .= True,+ "operation" .= operation,+ "queue" .= queue,+ "available" .= metrics.queueLength,+ "would_affect_at_most" .= affected+ ]+ }+ (forceInvocation env arguments)++handlePgmq :: IO (Either PgmqRuntimeError a) -> (a -> OpsOutcome) -> IO OpsOutcome+handlePgmq operation onSuccess = do+ result <- operation+ pure $ either (Failed . Text.pack . show) onSuccess result++pgmqRuntime :: OpsEnv -> JobRuntime+pgmqRuntime env = JobRuntime env.store.pool Nothing++rawJob :: Text -> Job Value+rawJob queue =+ Job+ { jobName = queue,+ jobQueue = queueRef queue,+ jobCodec = aesonJobCodec,+ jobPolicy = defaultRetryPolicy+ }++dlqListResult :: Text -> [DlqEntry Value] -> OpsResult+dlqListResult queue entries =+ OpsResult+ { headers = ["dlq_id", "queue", "original_id", "enqueued_at", "reads", "reason", "payload"],+ rows = map (dlqRow queue) entries,+ jsonValue = Aeson.toJSON (map (dlqJson queue) entries)+ }++dlqRow :: Text -> DlqEntry Value -> [Text]+dlqRow queue entry =+ [ showText entry.dlqMessageId,+ queue,+ maybe "" showText entry.originalMessageId,+ maybe "" (Text.pack . show) entry.originalEnqueuedAt,+ maybe "" showText entry.readCount,+ truncateCell 100 entry.reason,+ truncateCell 120 (jsonText (either (Aeson.String . Text.pack . show) id entry.originalPayload))+ ]++dlqJson :: Text -> DlqEntry Value -> Value+dlqJson queue entry =+ object+ [ "dlq_message_id" .= showText entry.dlqMessageId,+ "queue" .= queue,+ "reason" .= entry.reason,+ "original_payload" .= either (Aeson.String . Text.pack . show) id entry.originalPayload,+ "original_message_id" .= entry.originalMessageId,+ "original_enqueued_at" .= entry.originalEnqueuedAt,+ "read_count" .= entry.readCount,+ "raw_body" .= entry.rawBody+ ]++mutationCountResult :: Text -> Text -> Int -> OpsResult+mutationCountResult operation queue count =+ OpsResult+ { headers = ["operation", "queue", "affected"],+ rows = [[operation, queue, showText count]],+ jsonValue = object ["operation" .= operation, "queue" .= queue, "affected" .= count]+ }++showText :: (Show a) => a -> Text+showText = Text.pack . show++boolText :: Bool -> Text+boolText True = "true"+boolText False = "false"++forceInvocation :: OpsEnv -> [Text] -> Text+forceInvocation env arguments = Text.unwords (map shellQuote ("keiro-ops" : arguments <> globalFlags <> ["--force"]))+ where+ globalFlags = ["--json" | env.outputMode == Json] <> ["--allow-schema-drift" | env.allowSchemaDrift]++shellQuote :: Text -> Text+shellQuote value = "'" <> Text.replace "'" "'\"'\"'" value <> "'"
+ src/Keiro/Ops/Projection.hs view
@@ -0,0 +1,183 @@+-- | Operational adapters for projection dedup retention.+--+-- Commands call the public Keiro read-model and projection APIs in accordance+-- with ADR 28.+module Keiro.Ops.Projection+ ( Command (..),+ commandParser,+ isMutation,+ runCommand,+ )+where++import Data.Aeson (Value, object, (.=))+import Data.Int (Int64)+import Data.Text (Text)+import Data.Text qualified as Text+import Data.Time (UTCTime, defaultTimeLocale, formatTime, parseTimeM)+import Data.Vector qualified as Vector+import Effectful (Eff, IOE)+import Effectful.Error.Static (Error)+import Keiro.Ops.Env (OpsEnv (..), OutputMode (..))+import Keiro.Ops.Render+import Keiro.Projection (countAsyncProjectionDedupForBefore, pruneAsyncProjectionDedupForBefore)+import Keiro.ReadModel (storeHeadPosition)+import Kiroku.Store.Effect (Store, runStoreIO)+import Kiroku.Store.Error (StoreError)+import Kiroku.Store.Subscription+ ( SubscriptionCheckpoint (..),+ SubscriptionCheckpointInventory (..),+ SubscriptionName (..),+ subscriptionCheckpointInventory,+ )+import Kiroku.Store.Types (GlobalPosition (..))+import Options.Applicative hiding (action, value)+import Prelude++data Command+ = Position !Text+ | PruneDedup !Text !UTCTime+ deriving stock (Eq, Show)++commandParser :: Parser Command+commandParser =+ hsubparser+ ( command "position" (info positionParser (progDesc "Show one subscription's durable member checkpoints and floor"))+ <> command "prune-dedup" (info pruneParser (progDesc "Preview or prune one projection's old dedup rows"))+ )+ where+ positionParser = Position <$> textOption "subscription" "NAME" "Durable Kiroku subscription name"+ pruneParser =+ PruneDedup+ <$> textOption "projection" "NAME" "Async projection name"+ <*> option utcReader (long "before" <> metavar "UTC" <> help "Prune rows older than ISO-8601 UTC, for example 2026-08-01T00:00:00Z")++textOption :: String -> String -> String -> Parser Text+textOption name metavarText helpText = Text.pack <$> strOption (long name <> metavar metavarText <> help helpText)++utcReader :: ReadM UTCTime+utcReader = eitherReader $ \raw ->+ maybe+ (Left "expected ISO-8601 UTC such as 2026-08-01T00:00:00Z")+ Right+ (parseTimeM True defaultTimeLocale "%Y-%m-%dT%H:%M:%S%QZ" raw)++isMutation :: Command -> Bool+isMutation = \case+ Position {} -> False+ PruneDedup {} -> True++runCommand :: OpsEnv -> Command -> IO OpsOutcome+runCommand env = \case+ Position subscription -> runPosition env subscription+ PruneDedup projection before+ | env.force ->+ runAction env (pruneAsyncProjectionDedupForBefore projection before) $ \affected ->+ Succeeded (pruneResult False projection before affected)+ | otherwise ->+ runAction env (countAsyncProjectionDedupForBefore projection before) $ \affected ->+ PreviewRequired+ (pruneResult True projection before affected)+ (forceInvocation env ["projection", "prune-dedup", "--projection", projection, "--before", utcText before])++runPosition :: OpsEnv -> Text -> IO OpsOutcome+runPosition env subscription =+ runAction env action $ \(visibleHead, inventory) ->+ Succeeded (positionResult subscription visibleHead inventory)+ where+ action = (,) <$> storeHeadPosition <*> subscriptionCheckpointInventory++positionResult :: Text -> GlobalPosition -> SubscriptionCheckpointInventory -> OpsResult+positionResult requested visibleHead inventory =+ OpsResult+ { headers = ["subscription", "member", "checkpoint_position", "checkpoint_updated_at", "store_position", "visible_store_head", "global_position_distance", "minimum_checkpoint_position", "maximum_global_position_distance"],+ rows = humanRows,+ jsonValue =+ object+ [ "subscription" .= requested,+ "store_position" .= positionInt captured,+ "visible_store_head" .= positionInt visibleHead,+ "members" .= map (memberJson visibleHead) members,+ "minimum_checkpoint_position" .= fmap positionInt minimumCheckpoint,+ "maximum_global_position_distance" .= maximumDistance+ ]+ }+ where+ captured = storePosition inventory+ members =+ [ checkpoint+ | checkpoint@(SubscriptionCheckpoint (SubscriptionName name) _member _position _updatedAt) <-+ Vector.toList (checkpoints inventory),+ name == requested+ ]+ minimumCheckpoint = minimumMay [position | SubscriptionCheckpoint _ _ position _ <- members]+ maximumDistance = globalPositionDistance visibleHead <$> minimumCheckpoint+ summaryCells =+ [ maybe "" positionText minimumCheckpoint,+ maybe "" showText maximumDistance+ ]+ humanRows = case members of+ [] -> [[requested, "", "", "", positionText captured, positionText visibleHead, ""] <> summaryCells]+ _ -> map (memberRow captured visibleHead summaryCells) members++memberRow :: GlobalPosition -> GlobalPosition -> [Text] -> SubscriptionCheckpoint -> [Text]+memberRow captured visibleHead summaryCells (SubscriptionCheckpoint (SubscriptionName name) member position updatedAt) =+ [ name,+ showText member,+ positionText position,+ utcText updatedAt,+ positionText captured,+ positionText visibleHead,+ showText (globalPositionDistance visibleHead position)+ ]+ <> summaryCells++memberJson :: GlobalPosition -> SubscriptionCheckpoint -> Value+memberJson visibleHead (SubscriptionCheckpoint (SubscriptionName name) member position updatedAt) =+ object+ [ "subscription" .= name,+ "member" .= member,+ "checkpoint_position" .= positionInt position,+ "checkpoint_updated_at" .= updatedAt,+ "global_position_distance" .= globalPositionDistance visibleHead position+ ]++minimumMay :: (Ord a) => [a] -> Maybe a+minimumMay [] = Nothing+minimumMay values = Just (Prelude.minimum values)++globalPositionDistance :: GlobalPosition -> GlobalPosition -> Int64+globalPositionDistance (GlobalPosition captured) (GlobalPosition checkpoint) = max 0 (captured - checkpoint)++positionInt :: GlobalPosition -> Int64+positionInt (GlobalPosition value) = value++positionText :: GlobalPosition -> Text+positionText = showText . positionInt++runAction :: OpsEnv -> Eff '[Store, Error StoreError, IOE] a -> (a -> OpsOutcome) -> IO OpsOutcome+runAction env action onSuccess = do+ result <- runStoreIO env.store action+ pure $ either (Failed . Text.pack . show) onSuccess result++pruneResult :: Bool -> Text -> UTCTime -> Int64 -> OpsResult+pruneResult preview projection before affected =+ OpsResult+ { headers = ["projection", "before", if preview then "would_prune" else "pruned"],+ rows = [[projection, utcText before, showText affected]],+ jsonValue = object ["preview" .= preview, "projection" .= projection, "before" .= before, "affected" .= affected]+ }++utcText :: UTCTime -> Text+utcText = Text.pack . formatTime defaultTimeLocale "%Y-%m-%dT%H:%M:%S%QZ"++showText :: (Show a) => a -> Text+showText = Text.pack . show++forceInvocation :: OpsEnv -> [Text] -> Text+forceInvocation env arguments = Text.unwords (map shellQuote ("keiro-ops" : arguments <> globalFlags <> ["--force"]))+ where+ globalFlags = ["--json" | env.outputMode == Json] <> ["--allow-schema-drift" | env.allowSchemaDrift]++shellQuote :: Text -> Text+shellQuote value = "'" <> Text.replace "'" "'\"'\"'" value <> "'"
+ src/Keiro/Ops/Rebuild.hs view
@@ -0,0 +1,996 @@+-- | Catalog-backed projection rebuild commands.+--+-- Parsing, rendering, previews, and force policy live here; inventory and+-- rebuild semantics remain in 'ProjectionCatalogOperations'.+module Keiro.Ops.Rebuild+ ( Command (..),+ StartOptions (..),+ ResumeOptions (..),+ AbandonOptions (..),+ AdoptOptions (..),+ VersionedStartOptions (..),+ ReprojectStreamOptions (..),+ commandParser,+ isMutation,+ runCommand,+ streamReprojectionErrorCode,+ )+where++import Data.Aeson qualified as Aeson+import Data.Int (Int32, Int64)+import Data.List.NonEmpty (NonEmpty)+import Data.List.NonEmpty qualified as NonEmpty+import Data.Text (Text)+import Data.Text qualified as Text+import Data.Time (secondsToDiffTime)+import Data.UUID qualified as UUID+import Effectful (Eff, IOE)+import Effectful.Error.Static (Error)+import Keiro.Ops.Env (OpsEnv (..), OutputMode (..))+import Keiro.Ops.Parse (nonNegativeReader, readBoundedIntegral)+import Keiro.Ops.Render+import Keiro.Prelude ((&), (.~))+import Keiro.Projection.Catalog+ ( CatalogInventory (..),+ ExternalReadContractId,+ ExternalReadContractVersion (..),+ InventoryDedupKey (..),+ InventoryGroup (..),+ InventoryTarget (..),+ ProjectionId,+ ProjectionRevisionId,+ QualifiedTable (..),+ RebuildGroupId,+ StreamClearCount (..),+ TargetGenerationId (..),+ dedupKeyIdText,+ externalReadContractIdText,+ externalReadContractVersionValue,+ mkExternalReadContractId,+ mkProjectionId,+ mkProjectionRevisionId,+ mkRebuildGroupId,+ projectionIdText,+ projectionRevisionIdText,+ rebuildGroupIdText,+ targetIdText,+ )+import Keiro.Projection.Catalog.Operations+import Keiro.ReadModel.Rebuild+ ( GroupAdoptionClass (..),+ GroupLifecycleStatus (..),+ GroupRebuildMetadata (..),+ OrphanedRegistration (..),+ RebuildFailure (..),+ RebuildOptions (..),+ RebuildRequest (..),+ RebuildRunId,+ RebuildRunReport (..),+ RegistrationAdoption (..),+ RegistrationAdoptionAction (..),+ StreamReprojectionError (..),+ StreamReprojectionReport (..),+ StreamReprojectionRequest (..),+ VersionedRebuildReport (..),+ VersionedRetiredDropResult (..),+ VersionedRetiredGenerationPreview (..),+ VersionedTargetGeneration (..),+ VersionedTargetMode (..),+ defaultRebuildOptions,+ mkRebuildRunId,+ rebuildRunIdText,+ )+import Kiroku.Store.Effect (Store, runStoreIO)+import Kiroku.Store.Error (StoreError)+import Kiroku.Store.Types (GlobalPosition (..), StreamName (..), StreamVersion (..))+import Options.Applicative hiding (action, value)+import Options.Applicative qualified as Optparse++data Command+ = List+ | Preview !RebuildGroupId+ | Start !StartOptions+ | Status !RebuildRunId+ | Resume !ResumeOptions+ | Abandon !AbandonOptions+ | Adopt !AdoptOptions+ | VersionedStart !VersionedStartOptions+ | VersionedStatus !RebuildRunId+ | VersionedResume !RebuildRunId+ | VersionedAbandon !RebuildRunId+ | Retired+ | DropRetired !TargetGenerationId+ | ExternalReadStatus !ExternalReadContractId !ExternalReadContractVersion+ | RetireExternalRead !ExternalReadContractId !ExternalReadContractVersion+ | ReprojectStream !ReprojectStreamOptions+ deriving stock (Eq, Show)++data StartOptions = StartOptions+ { groupId :: !RebuildGroupId,+ runId :: !RebuildRunId,+ requestedBy :: !Text,+ reason :: !Text,+ replayFrom :: !GlobalPosition,+ pageSize :: !Int32+ }+ deriving stock (Eq, Show)++data ResumeOptions = ResumeOptions+ { runId :: !RebuildRunId,+ pageSize :: !Int32+ }+ deriving stock (Eq, Show)++data AbandonOptions = AbandonOptions+ { runId :: !RebuildRunId,+ failureCode :: !Text,+ failureDetail :: !Text+ }+ deriving stock (Eq, Show)++data AdoptOptions = AdoptOptions+ { groups :: !(NonEmpty RebuildGroupId)+ }+ deriving stock (Eq, Show)++data VersionedStartOptions = VersionedStartOptions+ { groupId :: !RebuildGroupId,+ runId :: !RebuildRunId,+ servingRevisionId :: !ProjectionRevisionId,+ candidateRevisionId :: !ProjectionRevisionId,+ targetMode :: !VersionedTargetMode,+ requestedBy :: !Text,+ reason :: !Text,+ pageSize :: !Int32,+ cutoverThreshold :: !Int64,+ cutoverLockTimeoutMs :: !Int64,+ promotionDedupLimit :: !Int64,+ retentionSeconds :: !Int64+ }+ deriving stock (Eq, Show)++data ReprojectStreamOptions = ReprojectStreamOptions+ { groupId :: !RebuildGroupId,+ projectionId :: !ProjectionId,+ streamName :: !StreamName,+ pageSize :: !Int32,+ maxEvents :: !Int64+ }+ deriving stock (Eq, Show)++commandParser :: Parser Command+commandParser =+ hsubparser+ ( command "list" (info (pure List) (progDesc "List the mounted catalog's rebuild groups"))+ <> command "preview" (info (Preview <$> groupArgument) (progDesc "Preview one catalog group without mutation"))+ <> command "start" (info (Start <$> startOptionsParser) (progDesc "Preview or start one catalog rebuild"))+ <> command "status" (info (Status <$> runArgument) (progDesc "Inspect one catalog rebuild run"))+ <> command "resume" (info (Resume <$> resumeOptionsParser) (progDesc "Preview or resume one catalog rebuild run"))+ <> command "abandon" (info (Abandon <$> abandonOptionsParser) (progDesc "Preview or abandon one catalog rebuild run"))+ <> command "adopt" (info (Adopt <$> adoptOptionsParser) (progDesc "Preview or adopt catalog slice changes for the named groups"))+ <> command "versioned" (info versionedCommandParser (progDesc "Operate schema-versioned online rebuilds"))+ <> command "retired" (info (pure Retired) (progDesc "List retired target generations"))+ <> command "drop-retired" (info (DropRetired <$> generationArgument) (progDesc "Preview or drop one retired target generation"))+ <> command "external-read" (info externalReadStatusParser (progDesc "Inspect one managed external read contract"))+ <> command "retire-external-read" (info retireExternalReadParser (progDesc "Preview or retire one managed external read contract"))+ <> command "reproject-stream" (info (ReprojectStream <$> reprojectStreamOptionsParser) (progDesc "Preview or repair one stream-scoped projection"))+ )++externalReadStatusParser :: Parser Command+externalReadStatusParser = ExternalReadStatus <$> externalReadContractArgument <*> externalReadVersionArgument++retireExternalReadParser :: Parser Command+retireExternalReadParser = RetireExternalRead <$> externalReadContractArgument <*> externalReadVersionArgument++reprojectStreamOptionsParser :: Parser ReprojectStreamOptions+reprojectStreamOptionsParser =+ ReprojectStreamOptions+ <$> groupArgument+ <*> argument projectionReader (metavar "PROJECTION")+ <*> (StreamName . Text.pack <$> strArgument (metavar "STREAM"))+ <*> option positiveInt32Reader (long "page-size" <> metavar "N" <> Optparse.value 500 <> showDefault <> help "Events fetched per stream page")+ <*> option positiveInt64Reader (long "max-events" <> metavar "N" <> Optparse.value 1000 <> showDefault <> help "Maximum locked stream event count admitted for one repair")++versionedCommandParser :: Parser Command+versionedCommandParser =+ hsubparser+ ( command "start" (info (VersionedStart <$> versionedStartOptionsParser) (progDesc "Preview or start an online versioned rebuild"))+ <> command "status" (info (VersionedStatus <$> runArgument) (progDesc "Inspect an online versioned rebuild"))+ <> command "resume" (info (VersionedResume <$> runArgument) (progDesc "Preview or advance one durable online rebuild phase"))+ <> command "abandon" (info (VersionedAbandon <$> runArgument) (progDesc "Preview or abandon an online versioned rebuild"))+ )++versionedStartOptionsParser :: Parser VersionedStartOptions+versionedStartOptionsParser =+ VersionedStartOptions+ <$> groupArgument+ <*> runIdOption+ <*> option revisionReader (long "serving-revision" <> metavar "REVISION")+ <*> option revisionReader (long "candidate-revision" <> metavar "REVISION")+ <*> option targetModeReader (long "target-mode" <> metavar "application|clone" <> Optparse.value ApplicationProvisioned <> showDefaultWith (const "application"))+ <*> textOption "requested-by" "IDENTITY" "Operator or automation identity"+ <*> textOption "reason" "TEXT" "Reason for this rebuild"+ <*> option positiveInt32Reader (long "page-size" <> metavar "N" <> Optparse.value 500 <> showDefault)+ <*> option nonNegativeInt64Reader (long "cutover-threshold" <> metavar "POSITIONS" <> Optparse.value 1000 <> showDefault)+ <*> option positiveInt64Reader (long "lock-timeout-ms" <> metavar "MILLISECONDS" <> Optparse.value 5000 <> showDefault)+ <*> option positiveInt64Reader (long "promotion-dedup-limit" <> metavar "ROWS" <> Optparse.value 1000000 <> showDefault)+ <*> option positiveInt64Reader (long "retention-seconds" <> metavar "SECONDS" <> Optparse.value 3600 <> showDefault)++startOptionsParser :: Parser StartOptions+startOptionsParser =+ StartOptions+ <$> groupArgument+ <*> runIdOption+ <*> textOption "requested-by" "IDENTITY" "Operator or automation identity"+ <*> textOption "reason" "TEXT" "Reason for this rebuild"+ <*> (GlobalPosition <$> option (nonNegativeReader "expected a non-negative global position") (long "from" <> metavar "POSITION" <> Optparse.value 0 <> showDefault <> help "Inclusive replay start position"))+ <*> option positiveInt32Reader (long "page-size" <> metavar "N" <> Optparse.value 500 <> showDefault <> help "Events fetched per replay page")++resumeOptionsParser :: Parser ResumeOptions+resumeOptionsParser =+ ResumeOptions+ <$> runArgument+ <*> option positiveInt32Reader (long "page-size" <> metavar "N" <> Optparse.value 500 <> showDefault <> help "Events fetched per replay page")++abandonOptionsParser :: Parser AbandonOptions+abandonOptionsParser =+ AbandonOptions+ <$> runArgument+ <*> textOption "code" "CODE" "Stable failure code"+ <*> textOption "detail" "TEXT" "Operator-visible failure detail"++adoptOptionsParser :: Parser AdoptOptions+adoptOptionsParser = AdoptOptions . NonEmpty.fromList <$> some groupArgument++groupArgument :: Parser RebuildGroupId+groupArgument = argument groupReader (metavar "GROUP")++runArgument :: Parser RebuildRunId+runArgument = argument runReader (metavar "RUN_ID")++runIdOption :: Parser RebuildRunId+runIdOption = option runReader (long "run-id" <> metavar "RUN_ID" <> help "Stable identity for this rebuild attempt")++generationArgument :: Parser TargetGenerationId+generationArgument = argument generationReader (metavar "GENERATION_ID")++externalReadContractArgument :: Parser ExternalReadContractId+externalReadContractArgument = argument externalReadContractReader (metavar "CONTRACT")++externalReadVersionArgument :: Parser ExternalReadContractVersion+externalReadVersionArgument = argument externalReadVersionReader (metavar "VERSION")++textOption :: String -> String -> String -> Parser Text+textOption name metavarText helpText = Text.pack <$> strOption (long name <> metavar metavarText <> help helpText)++groupReader :: ReadM RebuildGroupId+groupReader = eitherReader (firstShow . mkRebuildGroupId . Text.pack)++runReader :: ReadM RebuildRunId+runReader = eitherReader (firstText . mkRebuildRunId . Text.pack)++revisionReader :: ReadM ProjectionRevisionId+revisionReader = eitherReader (firstShow . mkProjectionRevisionId . Text.pack)++projectionReader :: ReadM ProjectionId+projectionReader = eitherReader (firstShow . mkProjectionId . Text.pack)++generationReader :: ReadM TargetGenerationId+generationReader = eitherReader $ \raw ->+ maybe (Left "expected a UUID target generation id") (Right . TargetGenerationId) (UUID.fromString raw)++externalReadContractReader :: ReadM ExternalReadContractId+externalReadContractReader = eitherReader (firstShow . mkExternalReadContractId . Text.pack)++externalReadVersionReader :: ReadM ExternalReadContractVersion+externalReadVersionReader = eitherReader $ \raw ->+ case readBoundedIntegral raw of+ Just value | value > 0 -> Right (ExternalReadContractVersion value)+ _ -> Left "expected a positive external read contract version"++targetModeReader :: ReadM VersionedTargetMode+targetModeReader = eitherReader $ \case+ "application" -> Right ApplicationProvisioned+ "clone" -> Right RestrictedClone+ _ -> Left "expected application or clone"++firstShow :: (Show err) => Either err value -> Either String value+firstShow = either (Left . show) Right++firstText :: Either Text value -> Either String value+firstText = either (Left . Text.unpack) Right++positiveInt32Reader :: ReadM Int32+positiveInt32Reader = eitherReader $ \raw ->+ case readBoundedIntegral raw of+ Just value | value > 0 -> Right value+ _ -> Left "expected a positive 32-bit integer"++positiveInt64Reader :: ReadM Int64+positiveInt64Reader = eitherReader $ \raw ->+ case readBoundedIntegral raw of+ Just value | value > 0 -> Right value+ _ -> Left "expected a positive 64-bit integer"++nonNegativeInt64Reader :: ReadM Int64+nonNegativeInt64Reader = eitherReader $ \raw ->+ case readBoundedIntegral raw of+ Just value | value >= 0 -> Right value+ _ -> Left "expected a non-negative 64-bit integer"++isMutation :: Command -> Bool+isMutation = \case+ List -> False+ Preview {} -> False+ Status {} -> False+ Start {} -> True+ Resume {} -> True+ Abandon {} -> True+ Adopt {} -> True+ VersionedStart {} -> True+ VersionedStatus {} -> False+ VersionedResume {} -> True+ VersionedAbandon {} -> True+ Retired -> False+ DropRetired {} -> True+ ExternalReadStatus {} -> False+ RetireExternalRead {} -> True+ ReprojectStream {} -> True++runCommand :: OpsEnv -> ProjectionCatalogOperations -> Command -> IO OpsOutcome+runCommand env operations = \case+ List -> pure (Succeeded (inventoryResult (catalogInventoryReport operations)))+ Preview groupId ->+ runCatalogAction env (previewRegisteredGroupRebuild operations groupId) (Succeeded . registeredPreviewResult)+ Start options+ | env.force ->+ runCatalogAction env (startGroupRebuild operations options.groupId (startRebuildOptions options)) (Succeeded . runResult)+ | otherwise ->+ runCatalogAction env (previewRegisteredGroupRebuild operations options.groupId) $ \preview ->+ PreviewRequired (registeredPreviewResult preview) (forceInvocation env (startArguments options))+ Status runId ->+ runCatalogAction env (inspectGroupRebuild operations runId) (Succeeded . runResult)+ Resume options+ | env.force ->+ runCatalogAction env (resumeGroupRebuild operations options.runId (resumeRebuildOptions options)) (Succeeded . runResult)+ | otherwise ->+ runCatalogAction env (inspectGroupRebuild operations options.runId) $ \report ->+ PreviewRequired (runResult report) (forceInvocation env (resumeArguments options))+ Abandon options+ | env.force ->+ runCatalogAction+ env+ (abandonGroupRebuild operations options.runId (RebuildFailure options.failureCode options.failureDetail))+ (Succeeded . runResult)+ | otherwise ->+ runCatalogAction env (inspectGroupRebuild operations options.runId) $ \report ->+ PreviewRequired (runResult report) (forceInvocation env (abandonArguments options))+ Adopt options+ | env.force ->+ runCatalogAction env (adoptCatalogGroups operations options.groups) (Succeeded . adoptionOutcomeResult)+ | otherwise ->+ runCatalogAction env (previewCatalogAdoption operations options.groups) $ \report ->+ PreviewRequired (adoptionPreviewResult report) (forceInvocation env (adoptArguments options))+ VersionedStart options+ | env.force ->+ runCatalogAction+ env+ (startVersionedGroupRebuild operations (catalogVersionedStartOptions options))+ (Succeeded . versionedRunResult)+ | otherwise ->+ runCatalogAction env (previewRegisteredGroupRebuild operations options.groupId) $ \preview ->+ PreviewRequired+ (registeredPreviewResult preview)+ (forceInvocation env (versionedStartArguments options))+ VersionedStatus runId ->+ runCatalogAction env (inspectVersionedGroupRebuild operations runId) (Succeeded . versionedRunResult)+ VersionedResume runId+ | env.force ->+ runCatalogAction env (resumeVersionedGroupRebuild operations runId) (Succeeded . versionedRunResult)+ | otherwise ->+ runCatalogAction env (inspectVersionedGroupRebuild operations runId) $ \report ->+ PreviewRequired (versionedRunResult report) (forceInvocation env (versionedResumeArguments runId))+ VersionedAbandon runId+ | env.force ->+ runCatalogAction env (abandonVersionedGroupRebuild operations runId) (Succeeded . versionedRunResult)+ | otherwise ->+ runCatalogAction env (inspectVersionedGroupRebuild operations runId) $ \report ->+ PreviewRequired (versionedRunResult report) (forceInvocation env (versionedAbandonArguments runId))+ Retired ->+ runCatalogValue env (listRetiredGenerations operations) (Succeeded . retiredGenerationsResult)+ DropRetired generationId+ | env.force ->+ runCatalogAction env (dropRetiredGeneration operations generationId) (Succeeded . retiredDropResult)+ | otherwise ->+ runCatalogAction env (previewRetiredGenerationDrop operations generationId) $ \report ->+ PreviewRequired+ (retiredDropResult report)+ (forceInvocation env (dropRetiredArguments generationId))+ ExternalReadStatus contractId contractVersion ->+ runCatalogAction+ env+ (inspectExternalReadContract operations contractId contractVersion)+ (Succeeded . externalReadRetirementResult)+ RetireExternalRead contractId contractVersion+ | env.force ->+ runCatalogAction+ env+ (retireExternalReadContract operations contractId contractVersion)+ (Succeeded . externalReadRetirementResult)+ | otherwise ->+ runCatalogAction env (inspectExternalReadContract operations contractId contractVersion) $ \report ->+ PreviewRequired+ (externalReadRetirementResult report)+ (forceInvocation env (retireExternalReadArguments contractId contractVersion))+ ReprojectStream options+ | env.force ->+ runCatalogAction+ env+ (reprojectCatalogStream operations (streamReprojectionRequest options))+ (Succeeded . streamReprojectionResult)+ | otherwise ->+ runCatalogAction env (previewStreamReprojection operations (streamReprojectionRequest options)) $ \report ->+ PreviewRequired+ (streamReprojectionPreviewResult report)+ (forceInvocation env (reprojectStreamArguments options))++streamReprojectionRequest :: ReprojectStreamOptions -> StreamReprojectionRequest+streamReprojectionRequest options =+ StreamReprojectionRequest+ { rebuildGroupId = options.groupId,+ projectionId = options.projectionId,+ streamName = options.streamName,+ pageSize = options.pageSize,+ maxEvents = options.maxEvents+ }++catalogVersionedStartOptions :: VersionedStartOptions -> CatalogVersionedStartOptions+catalogVersionedStartOptions options =+ CatalogVersionedStartOptions+ { rebuildRunId = options.runId,+ rebuildGroupId = options.groupId,+ servingRevisionId = options.servingRevisionId,+ candidateRevisionId = options.candidateRevisionId,+ targetMode = options.targetMode,+ replayPageSize = options.pageSize,+ cutoverThreshold = options.cutoverThreshold,+ cutoverLockTimeoutMs = options.cutoverLockTimeoutMs,+ promotionDedupLimit = options.promotionDedupLimit,+ retentionDuration = secondsToDiffTime (fromIntegral options.retentionSeconds),+ requestedBy = options.requestedBy,+ requestReason = options.reason+ }++startRebuildOptions :: StartOptions -> RebuildOptions+startRebuildOptions options =+ ( defaultRebuildOptions+ RebuildRequest+ { rebuildRunId = options.runId,+ requestedBy = options.requestedBy,+ requestReason = options.reason,+ replayFrom = options.replayFrom+ }+ )+ & #replayPageSize+ .~ options.pageSize++resumeRebuildOptions :: ResumeOptions -> RebuildOptions+resumeRebuildOptions options =+ ( defaultRebuildOptions+ RebuildRequest+ { rebuildRunId = options.runId,+ requestedBy = "keiro-ops",+ requestReason = "resume existing rebuild",+ replayFrom = GlobalPosition 0+ }+ )+ & #replayPageSize+ .~ options.pageSize++runCatalogAction ::+ OpsEnv ->+ Eff '[Store, Error StoreError, IOE] (Either CatalogOpsError value) ->+ (value -> OpsOutcome) ->+ IO OpsOutcome+runCatalogAction env action onSuccess = do+ result <- runStoreIO env.store action+ pure $ case result of+ Left storeError -> Failed (Text.pack (show storeError))+ Right (Left catalogError) -> Failed (catalogOpsErrorText catalogError)+ Right (Right value) -> onSuccess value++runCatalogValue ::+ OpsEnv ->+ Eff '[Store, Error StoreError, IOE] value ->+ (value -> OpsOutcome) ->+ IO OpsOutcome+runCatalogValue env action onSuccess = do+ result <- runStoreIO env.store action+ pure $ either (Failed . Text.pack . show) onSuccess result++inventoryResult :: CatalogInventoryReport -> OpsResult+inventoryResult report =+ OpsResult+ { headers = ["group", "targets", "verifications", "slice_fingerprint", "catalog_fingerprint"],+ rows =+ [ [ rebuildGroupIdText group.rebuildGroupId,+ showText (length group.orderedTargets),+ showText (length group.verifications),+ maybe "" id (lookup group.rebuildGroupId report.groupSlices),+ report.catalogFingerprint+ ]+ | group <- report.inventory.inventoryGroups+ ],+ jsonValue = Aeson.toJSON report+ }++registeredPreviewResult :: RegisteredRebuildPreview -> OpsResult+registeredPreviewResult report =+ OpsResult+ { headers = ["group", "targets", "subscriptions", "dedup_keys", "destructive", "registered", "slice_fingerprint", "registered_slice_matches"],+ rows =+ [ [ rebuildGroupIdText preview.rebuildGroupId,+ showText (length preview.targets),+ showText (length preview.subscriptionResets),+ showText (length preview.dedupResets),+ boolText preview.destructive,+ maybe "no" (const "yes") report.registeredState,+ preview.sliceFingerprint,+ maybe "" boolText report.registeredSliceMatches+ ]+ ],+ jsonValue = Aeson.toJSON report+ }+ where+ preview = report.preview++adoptionPreviewResult :: CatalogAdoptionReport -> OpsResult+adoptionPreviewResult report =+ OpsResult+ { headers = ["name", "kind", "state", "scope", "stored", "current"],+ rows =+ [ [ rebuildGroupIdText group.rebuildGroupId,+ "group",+ adoptionStateText group.classification,+ adoptionScopeText group.inScope,+ maybe "" id group.storedSlice,+ group.currentSlice+ ]+ | group <- report.groups+ ]+ <> [ [rebuildGroupIdText groupId, "group", "removed", "skip", "", ""]+ | groupId <- report.removedGroups+ ]+ <> [ [ registration.registryName,+ "registration",+ registrationActionText registration.action,+ adoptionScopeText registration.inScope,+ "",+ ""+ ]+ | registration <- report.registrations+ ]+ <> [ [ orphan.registryName,+ "registration",+ "orphaned-old-name",+ adoptionScopeText orphan.inScope,+ rebuildGroupIdText orphan.boundGroupId,+ ""+ ]+ | orphan <- report.orphanedRegistrations+ ]+ <> [["note", adoptionNote, "", "", "", ""]]+ <> warningRows,+ jsonValue = Aeson.toJSON report+ }+ where+ warningRows =+ [ [ "note",+ "warning: out-of-scope groups still drift and will fail startup registration until adopted: "+ <> renderGroupIds report.outOfScopeChangedGroups,+ "",+ "",+ "",+ ""+ ]+ | not (null report.outOfScopeChangedGroups)+ ]+ <> [ [ "note",+ "warning: requested groups not yet registered; --force will refuse: "+ <> renderGroupIds requestedNewGroups,+ "",+ "",+ "",+ ""+ ]+ | not (null requestedNewGroups)+ ]+ requestedNewGroups =+ [ group.rebuildGroupId+ | group <- report.groups,+ group.inScope,+ group.classification == AdoptionNew+ ]++adoptionOutcomeResult :: CatalogAdoptionOutcome -> OpsResult+adoptionOutcomeResult outcome =+ OpsResult+ { headers = ["name", "kind", "outcome", "detail"],+ rows =+ [ [ rebuildGroupIdText metadata.rebuildGroupId,+ "group",+ lifecycleStatusText metadata.status,+ metadata.sliceFingerprint+ ]+ | metadata <- outcome.adoptedGroups+ ]+ <> [ [ registration.registryName,+ "registration",+ registrationOutcomeText registration.action,+ rebuildGroupIdText registration.rebuildGroupId+ ]+ | registration <- outcome.registrationOutcomes+ ]+ <> [ [ orphan.registryName,+ "registration",+ "orphaned-old-name",+ rebuildGroupIdText orphan.boundGroupId+ ]+ | orphan <- outcome.removedOrphans+ ],+ jsonValue = Aeson.toJSON outcome+ }++adoptionScopeText :: Bool -> Text+adoptionScopeText True = "adopt"+adoptionScopeText False = "skip"++registrationActionText :: RegistrationAdoptionAction -> Text+registrationActionText = \case+ RegistrationUpdate -> "update"+ RegistrationInsert -> "insert"++registrationOutcomeText :: RegistrationAdoptionAction -> Text+registrationOutcomeText = \case+ RegistrationUpdate -> "adopted"+ RegistrationInsert -> "inserted"++renderGroupIds :: [RebuildGroupId] -> Text+renderGroupIds = Text.intercalate ", " . map rebuildGroupIdText++adoptionStateText :: GroupAdoptionClass -> Text+adoptionStateText = \case+ AdoptionNew -> "new"+ AdoptionUnchanged -> "unchanged"+ AdoptionSliceChanged {} -> "slice-changed"+ AdoptionStaleFormat {} -> "stale-format"++lifecycleStatusText :: GroupLifecycleStatus -> Text+lifecycleStatusText = \case+ GroupLive -> "live"+ GroupRebuilding -> "rebuilding"+ GroupFailed -> "failed"+ UnknownGroupStatus value -> value++adoptionNote :: Text+adoptionNote =+ "adoption changes only keiro-owned registration metadata; run 'rebuild start' if the change invalidates persisted rows"++runResult :: CatalogRunReport -> OpsResult+runResult report =+ OpsResult+ { headers = ["run", "group", "status", "group_slice", "captured_head", "sources", "adapters", "verifications"],+ rows =+ [ [ rebuildRunIdText run.rebuildRunId,+ rebuildGroupIdText run.rebuildGroupId,+ Text.pack (show run.runStatus),+ run.groupSliceFingerprint,+ globalPositionText run.capturedHead,+ showText (length run.sources),+ showText (length run.adapters),+ showText (length run.verifications)+ ]+ ],+ jsonValue = Aeson.toJSON report+ }+ where+ run = report.run++versionedRunResult :: CatalogVersionedRunReport -> OpsResult+versionedRunResult report =+ OpsResult+ { headers = ["run", "group", "phase", "serving_revision", "candidate_revision", "epoch", "captured_head", "dedup_limit", "staged_dedup", "dedup_provisional_head", "promotion_prepared", "sources", "serving_generations", "candidate_generations"],+ rows =+ [ [ rebuildRunIdText run.rebuildRunId,+ rebuildGroupIdText run.rebuildGroupId,+ showText run.phase,+ projectionRevisionIdText run.servingRevisionId,+ projectionRevisionIdText run.candidateRevisionId,+ showText run.servingEpoch,+ globalPositionText run.capturedHead,+ showText run.promotionDedupLimit,+ showText run.stagedDedupCount,+ maybe "" globalPositionText run.dedupProvisionalHead,+ boolText run.promotionPrepared,+ showText (length run.sources),+ showText (length run.servingGenerations),+ showText (length run.candidateGenerations)+ ]+ ],+ jsonValue = Aeson.toJSON report+ }+ where+ run = report.run++retiredGenerationsResult :: CatalogRetiredGenerationsReport -> OpsResult+retiredGenerationsResult report =+ OpsResult+ { headers = ["generation", "group", "target", "revision", "table", "oid", "lifecycle"],+ rows = map generationRow report.generations,+ jsonValue = Aeson.toJSON report+ }++retiredDropResult :: CatalogRetiredDropReport -> OpsResult+retiredDropResult report =+ OpsResult+ { headers = ["generation", "group", "target", "revision", "table", "droppable", "blockers"],+ rows =+ case report of+ CatalogRetiredDropPreview preview ->+ [ generationSummaryRow preview.generation+ <> [ boolText preview.droppable,+ Text.intercalate+ ","+ ( maybe [] (\runId -> ["active-run:" <> rebuildRunIdText runId]) preview.activeRunId+ <> map ("read-contract:" <>) preview.supportedReadContracts+ <> map ("postgres-dependency:" <>) preview.externalDependencies+ )+ ]+ ]+ CatalogRetiredDropOutcome outcome ->+ [ generationSummaryRow outcome.generation+ <> ["yes", if outcome.alreadyDropped then "already-dropped" else "dropped"]+ ],+ jsonValue = Aeson.toJSON report+ }++externalReadRetirementResult :: CatalogExternalReadRetirementReport -> OpsResult+externalReadRetirementResult report =+ OpsResult+ { headers = ["contract", "version", "function", "state", "surface_generation", "dependents", "execute_grants"],+ rows =+ [ [ externalReadContractIdText report.contractId,+ showText (externalReadContractVersionValue report.contractVersion),+ report.publicFunction,+ report.currentState,+ showText report.surfaceGeneration,+ Text.intercalate "," report.dependentObjects,+ Text.intercalate "," report.executeGrants+ ]+ ],+ jsonValue = Aeson.toJSON report+ }++streamReprojectionPreviewResult :: CatalogStreamReprojectionPreview -> OpsResult+streamReprojectionPreviewResult report =+ OpsResult+ { headers = ["group", "projection", "stream", "serving_revision", "targets", "dedup_keys", "stream_version", "event_count", "expected_dedup_claims", "max_events", "soft_deleted", "truncate_before", "eligible", "refusal"],+ rows =+ [ [ rebuildGroupIdText report.rebuildGroupId,+ projectionIdText report.projectionId,+ streamNameText report.streamName,+ projectionRevisionIdText report.servingRevisionId,+ Text.intercalate "," (map (\target -> targetIdText target.targetId) report.targets),+ Text.intercalate "," (map (\dedup -> dedupKeyIdText dedup.dedupKeyId) report.affectedDedup),+ maybe "" streamVersionText report.streamVersion,+ maybe "" showText report.eventCount,+ maybe "" showText report.expectedDedupClaims,+ showText report.maxEvents,+ boolText report.softDeleted,+ maybe "" streamVersionText report.truncateBefore,+ boolText report.eligible,+ maybe "" id report.refusal+ ]+ ],+ jsonValue = Aeson.toJSON report+ }++streamReprojectionResult :: CatalogStreamReprojectionReport -> OpsResult+streamReprojectionResult report =+ OpsResult+ { headers = ["group", "projection", "stream", "serving_revision", "stream_version", "max_events", "cleared_rows", "replayed", "applied", "dedup_inserted", "dedup_existing", "verified"],+ rows =+ [ [ rebuildGroupIdText repair.rebuildGroupId,+ projectionIdText repair.projectionId,+ streamNameText repair.streamName,+ projectionRevisionIdText repair.servingRevisionId,+ streamVersionText repair.streamVersion,+ showText repair.maxEvents,+ Text.intercalate "," [targetIdText count.targetId <> ":" <> showText count.clearedRows | count <- repair.clearedRows],+ showText repair.replayedEvents,+ showText repair.appliedEvents,+ showText repair.dedupInserted,+ showText repair.dedupExisting,+ boolText repair.verified+ ]+ ],+ jsonValue = Aeson.toJSON report+ }+ where+ repair = report.repair++generationRow :: VersionedTargetGeneration -> [Text]+generationRow generation =+ generationSummaryRow generation+ <> [showText generation.relationOid, showText generation.lifecycle]++generationSummaryRow :: VersionedTargetGeneration -> [Text]+generationSummaryRow generation =+ [ generationIdText generation.generationId,+ rebuildGroupIdText generation.rebuildGroupId,+ targetIdText generation.targetId,+ projectionRevisionIdText generation.revisionId,+ generation.physicalTable.schemaName <> "." <> generation.physicalTable.tableName+ ]++generationIdText :: TargetGenerationId -> Text+generationIdText (TargetGenerationId value) = UUID.toText value++startArguments :: StartOptions -> [Text]+startArguments options =+ [ "rebuild",+ "start",+ rebuildGroupIdText options.groupId,+ "--run-id",+ rebuildRunIdText options.runId,+ "--requested-by",+ options.requestedBy,+ "--reason",+ options.reason,+ "--from",+ globalPositionText options.replayFrom,+ "--page-size",+ showText options.pageSize+ ]++resumeArguments :: ResumeOptions -> [Text]+resumeArguments options =+ ["rebuild", "resume", rebuildRunIdText options.runId, "--page-size", showText options.pageSize]++abandonArguments :: AbandonOptions -> [Text]+abandonArguments options =+ [ "rebuild",+ "abandon",+ rebuildRunIdText options.runId,+ "--code",+ options.failureCode,+ "--detail",+ options.failureDetail+ ]++adoptArguments :: AdoptOptions -> [Text]+adoptArguments options =+ "rebuild" : "adopt" : map rebuildGroupIdText (NonEmpty.toList options.groups)++versionedStartArguments :: VersionedStartOptions -> [Text]+versionedStartArguments options =+ [ "rebuild",+ "versioned",+ "start",+ rebuildGroupIdText options.groupId,+ "--run-id",+ rebuildRunIdText options.runId,+ "--serving-revision",+ projectionRevisionIdText options.servingRevisionId,+ "--candidate-revision",+ projectionRevisionIdText options.candidateRevisionId,+ "--target-mode",+ case options.targetMode of+ ApplicationProvisioned -> "application"+ RestrictedClone -> "clone",+ "--requested-by",+ options.requestedBy,+ "--reason",+ options.reason,+ "--page-size",+ showText options.pageSize,+ "--cutover-threshold",+ showText options.cutoverThreshold,+ "--lock-timeout-ms",+ showText options.cutoverLockTimeoutMs,+ "--promotion-dedup-limit",+ showText options.promotionDedupLimit,+ "--retention-seconds",+ showText options.retentionSeconds+ ]++versionedResumeArguments :: RebuildRunId -> [Text]+versionedResumeArguments runId =+ ["rebuild", "versioned", "resume", rebuildRunIdText runId]++versionedAbandonArguments :: RebuildRunId -> [Text]+versionedAbandonArguments runId =+ ["rebuild", "versioned", "abandon", rebuildRunIdText runId]++dropRetiredArguments :: TargetGenerationId -> [Text]+dropRetiredArguments generationId =+ ["rebuild", "drop-retired", generationIdText generationId]++retireExternalReadArguments :: ExternalReadContractId -> ExternalReadContractVersion -> [Text]+retireExternalReadArguments contractId contractVersion =+ [ "rebuild",+ "retire-external-read",+ externalReadContractIdText contractId,+ showText (externalReadContractVersionValue contractVersion)+ ]++reprojectStreamArguments :: ReprojectStreamOptions -> [Text]+reprojectStreamArguments options =+ [ "rebuild",+ "reproject-stream",+ rebuildGroupIdText options.groupId,+ projectionIdText options.projectionId,+ streamNameText options.streamName,+ "--page-size",+ showText options.pageSize,+ "--max-events",+ showText options.maxEvents+ ]++forceInvocation :: OpsEnv -> [Text] -> Text+forceInvocation env arguments = Text.unwords (map shellQuote ("keiro-ops" : arguments <> globalFlags <> ["--force"]))+ where+ globalFlags = ["--json" | env.outputMode == Json] <> ["--allow-schema-drift" | env.allowSchemaDrift]++shellQuote :: Text -> Text+shellQuote value = "'" <> Text.replace "'" "'\"'\"'" value <> "'"++globalPositionText :: GlobalPosition -> Text+globalPositionText (GlobalPosition value) = showText value++streamNameText :: StreamName -> Text+streamNameText (StreamName value) = value++streamVersionText :: StreamVersion -> Text+streamVersionText (StreamVersion value) = showText value++catalogOpsErrorText :: CatalogOpsError -> Text+catalogOpsErrorText errorValue =+ case errorValue of+ CatalogOpsStreamReprojectionError streamError ->+ streamReprojectionErrorCode streamError <> ": " <> Text.pack (show streamError)+ _ -> Text.pack (show errorValue)++streamReprojectionErrorCode :: StreamReprojectionError -> Text+streamReprojectionErrorCode = \case+ StreamReprojectionInvalidPageSize {} -> "stream-reprojection-invalid-page-size"+ StreamReprojectionInvalidMaxEvents {} -> "stream-reprojection-invalid-max-events"+ StreamReprojectionEventLimitExceeded {} -> "stream-reprojection-event-limit-exceeded"+ StreamReprojectionGroupUnregistered {} -> "stream-reprojection-group-unregistered"+ StreamReprojectionActiveRebuild {} -> "stream-reprojection-active-rebuild"+ StreamReprojectionGroupUnavailable {} -> "stream-reprojection-group-unavailable"+ StreamReprojectionSliceDrift {} -> "stream-reprojection-slice-drift"+ StreamReprojectionServingRevisionUnavailable {} -> "stream-reprojection-serving-revision-unavailable"+ StreamReprojectionServingBindingInvalid {} -> "stream-reprojection-serving-binding-invalid"+ StreamReprojectionUnknownProjection {} -> "stream-reprojection-unknown-projection"+ StreamReprojectionProjectionGroupMismatch {} -> "stream-reprojection-projection-group-mismatch"+ StreamReprojectionPolicyUnavailable {} -> "stream-reprojection-policy-unavailable"+ StreamReprojectionSourceMismatch {} -> "stream-reprojection-source-mismatch"+ StreamReprojectionHistoryUnavailable {} -> "stream-reprojection-history-unavailable"+ StreamReprojectionSoftDeleted {} -> "stream-reprojection-soft-deleted"+ StreamReprojectionTruncated {} -> "stream-reprojection-truncated"+ StreamReprojectionForeignEvent {} -> "stream-reprojection-foreign-event"+ StreamReprojectionClearFailed {} -> "stream-reprojection-clear-failed"+ StreamReprojectionClearEvidenceInvalid {} -> "stream-reprojection-clear-evidence-invalid"+ StreamReprojectionDecodeFailed {} -> "stream-reprojection-decode-failed"+ StreamReprojectionVerificationFailed {} -> "stream-reprojection-verification-failed"+ StreamReprojectionDedupIdentityUnavailable {} -> "stream-reprojection-dedup-identity-unavailable"+ StreamReprojectionHistoryIncomplete {} -> "stream-reprojection-history-incomplete"++boolText :: Bool -> Text+boolText True = "yes"+boolText False = "no"++showText :: (Show value) => value -> Text+showText = Text.pack . show
+ src/Keiro/Ops/Render.hs view
@@ -0,0 +1,89 @@+module Keiro.Ops.Render+ ( OpsOutcome (..),+ OpsResult (..),+ emptyResult,+ jsonText,+ messageResult,+ renderHuman,+ renderResult,+ truncateCell,+ )+where++import Data.Aeson (Value, object, (.=))+import Data.Aeson qualified as Aeson+import Data.ByteString.Lazy qualified as LazyByteString.Raw+import Data.ByteString.Lazy.Char8 qualified as LazyByteString+import Data.List (transpose)+import Data.Text (Text)+import Data.Text qualified as Text+import Data.Text.Encoding qualified as Text.Encoding+import Data.Text.IO qualified as Text.IO+import Keiro.Ops.Env (OpsEnv (..), OutputMode (..))+import System.Exit (ExitCode)++data OpsResult = OpsResult+ { headers :: ![Text],+ rows :: ![[Text]],+ jsonValue :: !Value+ }+ deriving stock (Eq, Show)++data OpsOutcome+ = Succeeded !OpsResult+ | SucceededWithExit !OpsResult !ExitCode+ | PreviewRequired !OpsResult !Text+ | Failed !Text+ deriving stock (Eq, Show)++emptyResult :: OpsResult+emptyResult = OpsResult [] [] (Aeson.Array mempty)++messageResult :: Text -> OpsResult+messageResult message =+ OpsResult+ { headers = ["message"],+ rows = [[message]],+ jsonValue = object ["message" .= message]+ }++jsonText :: Value -> Text+jsonText = Text.Encoding.decodeUtf8 . LazyByteString.Raw.toStrict . Aeson.encode++truncateCell :: Int -> Text -> Text+truncateCell limit value+ | Text.length value <= limit = value+ | limit <= 1 = Text.take limit value+ | otherwise = Text.take (limit - 1) value <> "…"++renderResult :: OpsEnv -> OpsResult -> IO ()+renderResult env result =+ case env.outputMode of+ HumanTable -> Text.IO.putStrLn (renderHuman result)+ Json -> LazyByteString.putStrLn (Aeson.encode result.jsonValue)++renderHuman :: OpsResult -> Text+renderHuman OpsResult {headers, rows}+ | null headers = ""+ | otherwise =+ Text.unlines+ ( renderRow widths headers+ : renderSeparator widths+ : map (renderRow widths . normalizeRow (length headers)) rows+ )+ where+ normalizedRows = map (normalizeRow (length headers)) rows+ columns = transpose (headers : normalizedRows)+ widths = map (maximum . map Text.length) columns++normalizeRow :: Int -> [Text] -> [Text]+normalizeRow width row = take width (row <> repeat "")++renderRow :: [Int] -> [Text] -> Text+renderRow widths cells =+ Text.intercalate " " (zipWith pad widths cells)+ where+ pad width cell = cell <> Text.replicate (width - Text.length cell) " "++renderSeparator :: [Int] -> Text+renderSeparator = Text.intercalate " " . map (`Text.replicate` "-")
+ src/Keiro/Ops/ReplayAudit.hs view
@@ -0,0 +1,185 @@+{-# LANGUAGE GADTs #-}++module Keiro.Ops.ReplayAudit+ ( Command (..),+ OpsAuditConfig (..),+ commandParser,+ runCommand,+ )+where++import Data.Aeson (Value, object, (.=))+import Data.Set qualified as Set+import Data.Text (Text)+import Data.Text qualified as Text+import Keiro.Ops.Env (OpsEnv (..))+import Keiro.Ops.Parse (nonNegativeReader, positiveIntReader)+import Keiro.Ops.Render+import Keiro.ReplayAudit qualified as Audit+import Kiroku.Store.Effect (runStoreIO)+import Kiroku.Store.Types+ ( EventType (..),+ GlobalPosition (..),+ StreamName (..),+ StreamVersion (..),+ )+import Options.Applicative hiding (value)+import Options.Applicative qualified as Optparse+import System.Exit (ExitCode (..))++newtype OpsAuditConfig = OpsAuditConfig+ { targets :: [Audit.SomeAuditTarget]+ }++newtype Command = Audit AuditOptions+ deriving stock (Eq, Show)++data AuditOptions = AuditOptions+ { mode :: !Audit.AuditMode,+ category :: !(Maybe Text),+ maxStreams :: !(Maybe Int),+ parallelism :: !Int,+ resumeFrom :: !(Maybe GlobalPosition)+ }+ deriving stock (Eq, Show)++commandParser :: Parser Command+commandParser =+ Audit+ <$> ( AuditOptions+ <$> auditModeParser+ <*> optional (Text.pack <$> strOption (long "category" <> metavar "CATEGORY" <> help "Run only the configured audit target for this category"))+ <*> optional (option positiveIntReader (long "budget" <> metavar "STREAMS" <> help "Maximum streams to inspect in this invocation"))+ <*> option positiveIntReader (long "parallelism" <> metavar "N" <> Optparse.value 4 <> showDefault <> help "Maximum concurrent stream audits")+ <*> optional (GlobalPosition <$> option (nonNegativeReader "expected a non-negative global position") (long "resume-from" <> metavar "POSITION" <> help "Resume after this global-position checkpoint"))+ )++auditModeParser :: Parser Audit.AuditMode+auditModeParser =+ flag' Audit.AuditFull (long "full" <> help "Audit every stream in each configured category")+ <|> ( Audit.AuditTargeted+ <$> ( Audit.AffectedSet+ <$> (Set.fromList . map EventType <$> some (Text.pack <$> strOption (long "target" <> metavar "EVENT_TYPE" <> help "Affected event type from replay-impact analysis; repeat as needed")))+ <*> switch (long "include-snapshots" <> help "Include streams selected through snapshot event types")+ )+ )++runCommand :: OpsEnv -> OpsAuditConfig -> Command -> IO OpsOutcome+runCommand env config (Audit options) =+ case selectedTargets of+ [] -> pure (Failed (missingTargetMessage options.category))+ targets -> do+ result <- runStoreIO env.store (Audit.auditTargets options.mode budget targets)+ pure $ case result of+ Left storeError -> Failed (Text.pack (show storeError))+ Right reports -> auditOutcome reports+ where+ selectedTargets =+ case options.category of+ Nothing -> config.targets+ Just wanted -> filter ((== wanted) . configuredCategory) config.targets+ budget =+ Audit.AuditBudget+ { maxStreams = options.maxStreams,+ parallelism = options.parallelism,+ resumeFrom = options.resumeFrom+ }++configuredCategory :: Audit.SomeAuditTarget -> Text+configuredCategory (Audit.SomeAuditTarget target) = target.category++missingTargetMessage :: Maybe Text -> Text+missingTargetMessage = \case+ Nothing -> "replay audit hook has no configured targets"+ Just category -> "replay audit hook has no target for category " <> category++auditOutcome :: [Audit.AuditReport] -> OpsOutcome+auditOutcome reports =+ let result = auditResult reports+ in case Audit.auditExitCode reports of+ 0 -> Succeeded result+ code -> SucceededWithExit result (ExitFailure code)++auditResult :: [Audit.AuditReport] -> OpsResult+auditResult reports =+ OpsResult+ { headers = ["category", "mode", "selected", "skipped", "failures", "divergences", "checkpoint"],+ rows = map reportRow reports,+ jsonValue = toJson reports+ }++reportRow :: Audit.AuditReport -> [Text]+reportRow report =+ [ report.targetCategory,+ report.mode,+ showText report.streamsSelected,+ showText report.streamsSkipped,+ showText report.failures,+ showText report.divergences,+ maybe "-" globalPositionText report.checkpoint+ ]++toJson :: [Audit.AuditReport] -> Value+toJson reports =+ object+ [ "schema" .= ("keiro/replay-audit/v1" :: Text),+ "exit_code" .= Audit.auditExitCode reports,+ "reports" .= map reportJson reports+ ]++reportJson :: Audit.AuditReport -> Value+reportJson report =+ object+ [ "category" .= report.targetCategory,+ "mode" .= report.mode,+ "streams_selected" .= report.streamsSelected,+ "streams_skipped" .= report.streamsSkipped,+ "failures" .= report.failures,+ "divergences" .= report.divergences,+ "checkpoint" .= fmap globalPositionValue report.checkpoint,+ "rejected_streams" .= map streamNameText report.rejectedStreams,+ "results" .= map streamResultJson report.results+ ]++streamResultJson :: Audit.StreamAuditResult -> Value+streamResultJson result =+ object+ [ "stream" .= streamNameText result.streamName,+ "outcome" .= outcomeJson result.outcome+ ]++outcomeJson :: Audit.AuditOutcome -> Value+outcomeJson = \case+ Audit.ReplayOk streamVersion digest ->+ object+ [ "kind" .= ("ok" :: Text),+ "stream_version" .= streamVersionValue streamVersion,+ "digest" .= digest+ ]+ Audit.ReplayFailed commandError ->+ object+ [ "kind" .= ("failed" :: Text),+ "error" .= Text.pack (show commandError)+ ]+ Audit.SeedDivergence seedVersion seededDigest fullDigest ->+ object+ [ "kind" .= ("seed-divergence" :: Text),+ "seed_version" .= streamVersionValue seedVersion,+ "seeded_digest" .= seededDigest,+ "full_digest" .= fullDigest+ ]++streamNameText :: StreamName -> Text+streamNameText (StreamName value) = value++streamVersionValue :: StreamVersion -> Int+streamVersionValue (StreamVersion value) = fromIntegral value++globalPositionValue :: GlobalPosition -> Integer+globalPositionValue (GlobalPosition value) = fromIntegral value++globalPositionText :: GlobalPosition -> Text+globalPositionText = showText . globalPositionValue++showText :: (Show a) => a -> Text+showText = Text.pack . show
+ src/Keiro/Ops/Shard.hs view
@@ -0,0 +1,155 @@+-- | Operational adapters for sharded-subscription ownership.+--+-- Ownership reads and releases use 'Keiro.Subscription.Shard' exclusively, as+-- required by ADR 28.+module Keiro.Ops.Shard+ ( Command (..),+ commandParser,+ isMutation,+ runCommand,+ )+where++import Data.Aeson (object, (.=))+import Data.Set qualified as Set+import Data.Text (Text)+import Data.Text qualified as Text+import Data.Time (UTCTime, getCurrentTime)+import Data.UUID (UUID)+import Data.UUID qualified as UUID+import Effectful (Eff, IOE, (:>))+import Effectful.Error.Static (Error)+import Keiro.Ops.Env (OpsEnv (..), OutputMode (..))+import Keiro.Ops.Render+import Keiro.Subscription.Shard+import Kiroku.Store.Effect (Store, runStoreIO)+import Kiroku.Store.Error (StoreError)+import Kiroku.Store.Subscription.Types (SubscriptionName (..))+import Options.Applicative hiding (action, value)++data Command+ = Status !Text+ | Relinquish !Text !WorkerId+ deriving stock (Eq, Show)++commandParser :: Parser Command+commandParser =+ hsubparser+ ( command "status" (info statusParser (progDesc "Show shard counts and ownership"))+ <> command "relinquish" (info relinquishParser (progDesc "Preview or release every bucket owned by one worker"))+ )+ where+ statusParser = Status <$> subscriptionOption+ relinquishParser = Relinquish <$> subscriptionOption <*> (WorkerId <$> option uuidReader (long "worker" <> metavar "UUID" <> help "Worker id to release"))++subscriptionOption :: Parser Text+subscriptionOption = Text.pack <$> strOption (long "subscription" <> metavar "NAME" <> help "Sharded subscription name")++uuidReader :: ReadM UUID+uuidReader = eitherReader $ \raw -> maybe (Left "expected a UUID worker id") Right (UUID.fromString raw)++isMutation :: Command -> Bool+isMutation = \case+ Status {} -> False+ Relinquish {} -> True++runCommand :: OpsEnv -> Command -> IO OpsOutcome+runCommand env = \case+ Status name -> do+ now <- getCurrentTime+ runAction env (statusAction name) (Succeeded . uncurry (statusResult now name))+ Relinquish name worker -> runRelinquish env name worker++statusAction :: (Store :> es) => Text -> Eff es ([(Int, Maybe WorkerId, Maybe UTCTime)], [(Int, Int)])+statusAction name = do+ ownership <- ownershipSnapshotFor (SubscriptionName name)+ counts <- shardCountSnapshot (SubscriptionName name)+ pure (ownership, counts)++runRelinquish :: OpsEnv -> Text -> WorkerId -> IO OpsOutcome+runRelinquish env name worker+ | not env.force =+ runAction env (ownershipSnapshotFor subscription) $ \ownership ->+ let buckets = ownedBuckets worker ownership+ in PreviewRequired+ (relinquishResult True name worker buckets)+ (forceInvocation env ["shard", "relinquish", "--subscription", name, "--worker", workerText worker])+ | otherwise =+ runAction env action $ \buckets ->+ Succeeded (relinquishResult False name worker buckets)+ where+ subscription = SubscriptionName name+ lease = ShardLease subscription worker 0 0+ action = do+ ownership <- ownershipSnapshotFor subscription+ let buckets = ownedBuckets worker ownership+ relinquish lease (Set.fromList buckets)+ pure buckets++ownedBuckets :: WorkerId -> [(Int, Maybe WorkerId, Maybe UTCTime)] -> [Int]+ownedBuckets worker rows = [bucket | (bucket, Just owner, _) <- rows, owner == worker]++runAction :: OpsEnv -> Eff '[Store, Error StoreError, IOE] a -> (a -> OpsOutcome) -> IO OpsOutcome+runAction env action onSuccess = do+ result <- runStoreIO env.store action+ pure $ either (Failed . Text.pack . show) onSuccess result++statusResult :: UTCTime -> Text -> [(Int, Maybe WorkerId, Maybe UTCTime)] -> [(Int, Int)] -> OpsResult+statusResult now name ownership counts =+ OpsResult+ { headers = ["subscription", "bucket", "owner", "lease_expires_at", "lease_state", "configured_shards"],+ rows = map rowText ownership,+ jsonValue = object ["subscription" .= name, "shard_counts" .= map countJson counts, "ownership" .= map ownershipJson ownership]+ }+ where+ configured = Text.intercalate "," [showText shardCount <> " (" <> showText rowCount <> " rows)" | (shardCount, rowCount) <- counts]+ rowText (bucket, owner, expiresAt) =+ [ name,+ showText bucket,+ maybe "unowned" workerText owner,+ maybe "" timeText expiresAt,+ leaseState now owner expiresAt,+ configured+ ]+ countJson (shardCount, rowCount) = object ["shard_count" .= shardCount, "rows" .= rowCount]+ ownershipJson (bucket, owner, expiresAt) =+ object+ [ "bucket" .= bucket,+ "owner" .= fmap workerText owner,+ "lease_expires_at" .= expiresAt,+ "lease_state" .= leaseState now owner expiresAt+ ]++leaseState :: UTCTime -> Maybe WorkerId -> Maybe UTCTime -> Text+leaseState _ Nothing _ = "unowned"+leaseState now (Just _) (Just expiry)+ | expiry < now = "expired"+ | otherwise = "live"+leaseState _ (Just _) Nothing = "invalid"++relinquishResult :: Bool -> Text -> WorkerId -> [Int] -> OpsResult+relinquishResult preview name worker buckets =+ OpsResult+ { headers = ["subscription", "worker", "bucket", "disposition"],+ rows = [[name, workerText worker, showText bucket, disposition] | bucket <- buckets],+ jsonValue = object ["preview" .= preview, "subscription" .= name, "worker" .= workerText worker, "buckets" .= buckets, "affected" .= length buckets]+ }+ where+ disposition = if preview then "would_relinquish" else "relinquished"++workerText :: WorkerId -> Text+workerText (WorkerId value) = UUID.toText value++timeText :: UTCTime -> Text+timeText = Text.pack . show++showText :: (Show a) => a -> Text+showText = Text.pack . show++forceInvocation :: OpsEnv -> [Text] -> Text+forceInvocation env arguments = Text.unwords (map shellQuote ("keiro-ops" : arguments <> globalFlags <> ["--force"]))+ where+ globalFlags = ["--json" | env.outputMode == Json] <> ["--allow-schema-drift" | env.allowSchemaDrift]++shellQuote :: Text -> Text+shellQuote value = "'" <> Text.replace "'" "'\"'\"'" value <> "'"
+ src/Keiro/Ops/Snapshot.hs view
@@ -0,0 +1,291 @@+-- | Operational adapters for Keiro's advisory snapshot cache.+--+-- Inspection and deletion use the public snapshot storage operations added to+-- the owning Keiro library, preserving ADR 3 and ADR 28.+module Keiro.Ops.Snapshot+ ( Command (..),+ ExpectedDiscriminators (..),+ PreflightEvidence (..),+ commandParser,+ isMutation,+ preflightFor,+ runCommand,+ )+where++import Data.Aeson (Value, object, (.=))+import Data.Aeson qualified as Aeson+import Data.ByteString.Lazy qualified as LazyByteString+import Data.Int (Int64)+import Data.Text (Text)+import Data.Text qualified as Text+import Effectful (Eff, IOE, (:>))+import Effectful.Error.Static (Error)+import Keiro.Ops.Env (OpsEnv (..), OutputMode (..))+import Keiro.Ops.Parse (nonNegativeIntReader, nonNegativeReader)+import Keiro.Ops.Render+import Keiro.Snapshot.Schema+import Keiro.Workflow.Snapshot (workflowStateCodecVersion, workflowStateShapeHash)+import Kiroku.Store.Effect (Store, runStoreIO)+import Kiroku.Store.Error (StoreError)+import Kiroku.Store.Read (getStream)+import Kiroku.Store.Types (StreamInfo (..), StreamName (..), StreamVersion (..))+import Options.Applicative hiding (action, info, value)+import Options.Applicative qualified as Opt++data ExpectedDiscriminators = ExpectedDiscriminators+ { stateCodecVersion :: !Int,+ regfileShapeHash :: !Text,+ stateShapeHash :: !Text+ }+ deriving stock (Eq, Show)++data Command+ = Show !Text+ | Delete !Text+ | TruncationPreflight !Text !StreamVersion !(Maybe ExpectedDiscriminators)+ deriving stock (Eq, Show)++data PreflightEvidence = PreflightEvidence+ { streamName :: !Text,+ truncateBefore :: !StreamVersion,+ requiredSnapshotVersion :: !StreamVersion,+ expectedDiscriminators :: !(Maybe ExpectedDiscriminators),+ snapshotRow :: !(Maybe SnapshotRow),+ versionCovered :: !Bool,+ discriminatorsMatch :: !Bool,+ passed :: !Bool,+ reason :: !Text+ }+ deriving stock (Eq, Show)++commandParser :: Parser Command+commandParser =+ hsubparser+ ( command "show" (Opt.info (Show <$> streamOption) (progDesc "Inspect the advisory snapshot row for a stream"))+ <> command "delete" (Opt.info (Delete <$> streamOption) (progDesc "Preview or delete a stream's advisory snapshot"))+ <> command "truncation-preflight" (Opt.info preflightParser (progDesc "Check snapshot coverage before moving a Kiroku truncate marker"))+ )+ where+ preflightParser =+ TruncationPreflight+ <$> streamOption+ <*> (StreamVersion <$> option (nonNegativeReader "expected a non-negative stream version") (long "before" <> metavar "VERSION" <> help "Proposed Kiroku truncate-before version"))+ <*> optional expectedParser+ expectedParser =+ ExpectedDiscriminators+ <$> option nonNegativeIntReader (long "state-codec-version" <> metavar "N" <> help "Application's current state codec version")+ <*> textOption "regfile-shape-hash" "HASH" "Application's current register-layout hash"+ <*> textOption "state-shape-hash" "HASH" "Application's current control-state/fold hash"++streamOption :: Parser Text+streamOption = textOption "stream" "NAME" "Kiroku stream name"++textOption :: String -> String -> String -> Parser Text+textOption name metavarText helpText = Text.pack <$> strOption (long name <> metavar metavarText <> help helpText)++isMutation :: Command -> Bool+isMutation = \case+ Show {} -> False+ Delete {} -> True+ TruncationPreflight {} -> False++runCommand :: OpsEnv -> Command -> IO OpsOutcome+runCommand env = \case+ Show name -> runAction env (lookupByName name) (Succeeded . snapshotResult name)+ Delete name -> runDelete env name+ TruncationPreflight name before suppliedExpected ->+ runAction env (preflightFor name before suppliedExpected) (Succeeded . preflightResult)++lookupByName :: (Store :> es) => Text -> Eff es (Maybe SnapshotRow)+lookupByName name = do+ stream <- getStream (StreamName name)+ case stream of+ Nothing -> pure Nothing+ Just info -> lookupSnapshotRow info.id++runDelete :: OpsEnv -> Text -> IO OpsOutcome+runDelete env name+ | not env.force =+ runAction env (lookupByName name) $ \row ->+ PreviewRequired+ (snapshotDeleteResult True name row False)+ (forceInvocation env ["snapshot", "delete", "--stream", name])+ | otherwise =+ runAction env action $ \(row, deleted) ->+ Succeeded (snapshotDeleteResult False name row deleted)+ where+ action = do+ stream <- getStream (StreamName name)+ case stream of+ Nothing -> pure (Nothing, False)+ Just info -> do+ row <- lookupSnapshotRow info.id+ deleted <- deleteSnapshotRow info.id+ pure (row, deleted)++-- | Evaluate the database-only part of the truncation guard. Workflow journal+-- streams have public fixed discriminators and are recognized automatically.+-- Aggregate streams require the application-owned current discriminator tuple+-- to be supplied explicitly; a standalone binary cannot infer compiled codecs.+preflightFor ::+ (Store :> es) =>+ Text ->+ StreamVersion ->+ Maybe ExpectedDiscriminators ->+ Eff es PreflightEvidence+preflightFor name before suppliedExpected = do+ row <- lookupByName name+ let expected = suppliedExpected <|> workflowExpected name+ required = predecessor before+ covered = maybe False ((>= required) . (.streamVersion)) row+ matches = case (expected, row) of+ (Just wanted, Just found) -> discriminatorMatches wanted found+ _ -> False+ (ok, explanation) = case (row, expected, covered, matches) of+ (Nothing, _, _, _) -> (False, "no snapshot row exists")+ (Just _, Nothing, _, _) -> (False, "current codec discriminators are required for a non-workflow stream")+ (Just _, Just _, False, _) -> (False, "snapshot version does not cover the proposed truncation boundary")+ (Just _, Just _, True, False) -> (False, "snapshot discriminators do not match the expected current codec")+ (Just _, Just _, True, True) -> (True, "snapshot covers the boundary and matches the expected codec")+ pure+ PreflightEvidence+ { streamName = name,+ truncateBefore = before,+ requiredSnapshotVersion = required,+ expectedDiscriminators = expected,+ snapshotRow = row,+ versionCovered = covered,+ discriminatorsMatch = matches,+ passed = ok,+ reason = explanation+ }++workflowExpected :: Text -> Maybe ExpectedDiscriminators+workflowExpected name+ | "wf:" `Text.isPrefixOf` name =+ Just+ ExpectedDiscriminators+ { stateCodecVersion = workflowStateCodecVersion,+ regfileShapeHash = workflowStateShapeHash,+ stateShapeHash = workflowStateShapeHash+ }+ | otherwise = Nothing++predecessor :: StreamVersion -> StreamVersion+predecessor (StreamVersion before) = StreamVersion (max 0 (before - 1))++discriminatorMatches :: ExpectedDiscriminators -> SnapshotRow -> Bool+discriminatorMatches expected row =+ expected.stateCodecVersion == row.stateCodecVersion+ && expected.regfileShapeHash == row.regfileShapeHash+ && expected.stateShapeHash == row.stateShapeHash++runAction :: OpsEnv -> Eff '[Store, Error StoreError, IOE] a -> (a -> OpsOutcome) -> IO OpsOutcome+runAction env action onSuccess = do+ result <- runStoreIO env.store action+ pure $ either (Failed . Text.pack . show) onSuccess result++snapshotResult :: Text -> Maybe SnapshotRow -> OpsResult+snapshotResult name row =+ OpsResult+ { headers = ["stream", "snapshot_version", "state_codec_version", "regfile_shape_hash", "state_shape_hash", "state_bytes", "updated_at"],+ rows = maybe [] (pure . snapshotRowText name) row,+ jsonValue = maybe Aeson.Null (snapshotJson name) row+ }++snapshotRowText :: Text -> SnapshotRow -> [Text]+snapshotRowText name row =+ [ name,+ versionText row.streamVersion,+ showText row.stateCodecVersion,+ row.regfileShapeHash,+ row.stateShapeHash,+ showText (LazyByteString.length (Aeson.encode row.state)),+ showText row.updatedAt+ ]++snapshotJson :: Text -> SnapshotRow -> Value+snapshotJson name row =+ object+ [ "stream" .= name,+ "stream_version" .= versionInt row.streamVersion,+ "state" .= row.state,+ "state_codec_version" .= row.stateCodecVersion,+ "regfile_shape_hash" .= row.regfileShapeHash,+ "state_shape_hash" .= row.stateShapeHash,+ "state_bytes" .= LazyByteString.length (Aeson.encode row.state),+ "created_at" .= row.createdAt,+ "updated_at" .= row.updatedAt+ ]++snapshotDeleteResult :: Bool -> Text -> Maybe SnapshotRow -> Bool -> OpsResult+snapshotDeleteResult preview name row deleted =+ OpsResult+ { headers = ["stream", "snapshot_version", "disposition"],+ rows = [[name, maybe "" (versionText . (.streamVersion)) row, disposition]],+ jsonValue = object ["preview" .= preview, "stream" .= name, "snapshot" .= fmap (snapshotJson name) row, "deleted" .= deleted, "disposition" .= disposition]+ }+ where+ disposition+ | preview = maybe "not_found" (const "would_delete") row+ | deleted = "deleted"+ | otherwise = "not_found"++preflightResult :: PreflightEvidence -> OpsResult+preflightResult evidence =+ OpsResult+ { headers = ["stream", "before", "required_snapshot", "snapshot_version", "version_covered", "discriminators_match", "passed", "reason"],+ rows =+ [ [ evidence.streamName,+ versionText evidence.truncateBefore,+ versionText evidence.requiredSnapshotVersion,+ maybe "none" (versionText . (.streamVersion)) evidence.snapshotRow,+ boolText evidence.versionCovered,+ boolText evidence.discriminatorsMatch,+ boolText evidence.passed,+ evidence.reason+ ]+ ],+ jsonValue =+ object+ [ "stream" .= evidence.streamName,+ "truncate_before" .= versionInt evidence.truncateBefore,+ "required_snapshot_version" .= versionInt evidence.requiredSnapshotVersion,+ "expected_discriminators" .= fmap expectedJson evidence.expectedDiscriminators,+ "snapshot" .= fmap (snapshotJson evidence.streamName) evidence.snapshotRow,+ "version_covered" .= evidence.versionCovered,+ "discriminators_match" .= evidence.discriminatorsMatch,+ "passed" .= evidence.passed,+ "reason" .= evidence.reason+ ]+ }++expectedJson :: ExpectedDiscriminators -> Value+expectedJson expected =+ object+ [ "state_codec_version" .= expected.stateCodecVersion,+ "regfile_shape_hash" .= expected.regfileShapeHash,+ "state_shape_hash" .= expected.stateShapeHash+ ]++versionInt :: StreamVersion -> Int64+versionInt (StreamVersion value) = value++versionText :: StreamVersion -> Text+versionText = showText . versionInt++boolText :: Bool -> Text+boolText True = "true"+boolText False = "false"++showText :: (Show a) => a -> Text+showText = Text.pack . show++forceInvocation :: OpsEnv -> [Text] -> Text+forceInvocation env arguments = Text.unwords (map shellQuote ("keiro-ops" : arguments <> globalFlags <> ["--force"]))+ where+ globalFlags = ["--json" | env.outputMode == Json] <> ["--allow-schema-drift" | env.allowSchemaDrift]++shellQuote :: Text -> Text+shellQuote value = "'" <> Text.replace "'" "'\"'\"'" value <> "'"
+ src/Keiro/Ops/Stream.hs view
@@ -0,0 +1,465 @@+-- | Operational adapters for Kiroku stream reads and lifecycle operations.+--+-- Every database action is a public @kiroku-store@ operation, preserving that+-- library's schema ownership under ADR 28.+module Keiro.Ops.Stream+ ( Command (..),+ TruncateCommand (..),+ commandParser,+ isMutation,+ runCommand,+ )+where++import Data.Aeson (Value, object, (.=))+import Data.Aeson qualified as Aeson+import Data.Int (Int64)+import Data.Map.Strict (Map)+import Data.Map.Strict qualified as Map+import Data.Text (Text)+import Data.Text qualified as Text+import Data.Text.IO qualified as Text.IO+import Data.Time (UTCTime, defaultTimeLocale, formatTime)+import Data.UUID (UUID)+import Data.UUID qualified as UUID+import Data.Vector qualified as Vector+import Effectful (Eff, IOE)+import Effectful.Error.Static (Error)+import Keiro.Ops.Env (OpsEnv (..), OutputMode (..))+import Keiro.Ops.Parse (nonNegativeIntReader, nonNegativeReader, positiveIntReader)+import Keiro.Ops.Render+import Keiro.Ops.Snapshot qualified as Snapshot+import Keiro.ReadModel (storeHeadPosition)+import Kiroku.Store.Causation (findCausationAncestors, findCausationDescendants)+import Kiroku.Store.Effect (Store, runStoreIO)+import Kiroku.Store.Error (StoreError)+import Kiroku.Store.Lifecycle+import Kiroku.Store.Read (getStream, lookupStreamNames, readStreamForward)+import Kiroku.Store.Subscription+ ( SubscriptionCheckpoint (..),+ SubscriptionCheckpointInventory (..),+ SubscriptionName (..),+ subscriptionCheckpointInventory,+ )+import Kiroku.Store.Types+import Options.Applicative hiding (action, info, value)+import Options.Applicative qualified as Opt+import System.IO (hFlush, stdout)++data Command+ = Show !Text !StreamVersion !Int+ | SoftDelete !Text+ | Undelete !Text+ | HardDelete !Text+ | TruncateBefore !TruncateCommand+ | Causation !EventId+ | Subscriptions+ deriving stock (Eq, Show)++data TruncateCommand+ = SetTruncateBefore !Text !StreamVersion !(Maybe Snapshot.ExpectedDiscriminators) !Bool+ | ClearTruncateBefore !Text+ deriving stock (Eq, Show)++commandParser :: Parser Command+commandParser =+ hsubparser+ ( command "show" (Opt.info showParser (progDesc "Show stream metadata and ordered events"))+ <> command "soft-delete" (Opt.info (SoftDelete <$> streamArgument) (progDesc "Preview or soft-delete a stream"))+ <> command "undelete" (Opt.info (Undelete <$> streamArgument) (progDesc "Preview or restore a soft-deleted stream"))+ <> command "hard-delete" (Opt.info (HardDelete <$> streamArgument) (progDesc "Preview or permanently delete a stream"))+ <> command "truncate-before" (Opt.info (TruncateBefore <$> truncateParser) (progDesc "Operate the reversible stream visibility marker"))+ <> command "causation" (Opt.info (Causation <$> eventIdArgument) (progDesc "Show an event's causation ancestors and descendants"))+ <> command "subscriptions" (Opt.info (pure Subscriptions) (progDesc "List durable subscription checkpoints"))+ )+ where+ showParser =+ Show+ <$> streamArgument+ <*> (StreamVersion <$> option (nonNegativeReader "expected a non-negative stream version") (long "from" <> metavar "VERSION" <> Opt.value 0 <> showDefault <> help "Exclusive stream-version cursor"))+ <*> option positiveIntReader (long "limit" <> metavar "N" <> Opt.value 100 <> showDefault <> help "Maximum events")++truncateParser :: Parser TruncateCommand+truncateParser =+ hsubparser+ ( command "set" (Opt.info setParser (progDesc "Preview or set a truncate-before marker after snapshot preflight"))+ <> command "clear" (Opt.info (ClearTruncateBefore <$> streamArgument) (progDesc "Preview or clear a truncate-before marker"))+ )+ where+ setParser =+ SetTruncateBefore+ <$> streamArgument+ <*> (StreamVersion <$> argument (nonNegativeReader "expected a non-negative stream version") (metavar "VERSION"))+ <*> optional expectedParser+ <*> switch (long "skip-preflight" <> help "Bypass snapshot coverage checking (dangerous)")+ expectedParser =+ Snapshot.ExpectedDiscriminators+ <$> option nonNegativeIntReader (long "state-codec-version" <> metavar "N" <> help "Application's current state codec version")+ <*> textOption "regfile-shape-hash" "HASH" "Application's current register-layout hash"+ <*> textOption "state-shape-hash" "HASH" "Application's current control-state/fold hash"++streamArgument :: Parser Text+streamArgument = Text.pack <$> argument str (metavar "STREAM")++textOption :: String -> String -> String -> Parser Text+textOption name metavarText helpText = Text.pack <$> strOption (long name <> metavar metavarText <> help helpText)++eventIdArgument :: Parser EventId+eventIdArgument = EventId <$> argument uuidReader (metavar "EVENT_ID")++uuidReader :: ReadM UUID+uuidReader = eitherReader $ \raw -> maybe (Left "expected a UUID event id") Right (UUID.fromString raw)++isMutation :: Command -> Bool+isMutation = \case+ Show {} -> False+ Causation {} -> False+ Subscriptions -> False+ SoftDelete {} -> True+ Undelete {} -> True+ HardDelete {} -> True+ TruncateBefore {} -> True++runCommand :: OpsEnv -> Command -> IO OpsOutcome+runCommand env = \case+ Show name from limit -> runShow env name from limit+ SoftDelete name -> runLifecycle env "soft-delete" name (softDeleteStream (StreamName name))+ Undelete name -> runLifecycle env "undelete" name (undeleteStream (StreamName name))+ HardDelete name -> runHardDelete env name+ TruncateBefore truncateCommand -> runTruncate env truncateCommand+ Causation eventId -> runCausation env eventId+ Subscriptions -> runSubscriptions env++runSubscriptions :: OpsEnv -> IO OpsOutcome+runSubscriptions env =+ runAction env action $ \(visibleHead, inventory) ->+ Succeeded (subscriptionInventoryResult visibleHead inventory)+ where+ action = (,) <$> storeHeadPosition <*> subscriptionCheckpointInventory++subscriptionInventoryResult :: GlobalPosition -> SubscriptionCheckpointInventory -> OpsResult+subscriptionInventoryResult visibleHead inventory =+ OpsResult+ { headers = ["subscription", "member", "checkpoint_position", "checkpoint_updated_at", "store_position", "visible_store_head", "global_position_distance"],+ rows = map (checkpointRow captured visibleHead) durableCheckpoints,+ jsonValue =+ object+ [ "store_position" .= positionInt captured,+ "visible_store_head" .= positionInt visibleHead,+ "checkpoints" .= map (checkpointJson visibleHead) durableCheckpoints+ ]+ }+ where+ captured = storePosition inventory+ durableCheckpoints = Vector.toList (checkpoints inventory)++checkpointRow :: GlobalPosition -> GlobalPosition -> SubscriptionCheckpoint -> [Text]+checkpointRow captured visibleHead (SubscriptionCheckpoint (SubscriptionName name) member position updatedAt) =+ [ name,+ showText member,+ positionText position,+ utcText updatedAt,+ positionText captured,+ positionText visibleHead,+ showText (globalPositionDistance visibleHead position)+ ]++checkpointJson :: GlobalPosition -> SubscriptionCheckpoint -> Value+checkpointJson visibleHead (SubscriptionCheckpoint (SubscriptionName name) member position updatedAt) =+ object+ [ "subscription" .= name,+ "member" .= member,+ "checkpoint_position" .= positionInt position,+ "checkpoint_updated_at" .= updatedAt,+ "global_position_distance" .= globalPositionDistance visibleHead position+ ]++globalPositionDistance :: GlobalPosition -> GlobalPosition -> Int64+globalPositionDistance (GlobalPosition captured) (GlobalPosition checkpoint) = max 0 (captured - checkpoint)++utcText :: UTCTime -> Text+utcText = Text.pack . formatTime defaultTimeLocale "%Y-%m-%dT%H:%M:%S%QZ"++runShow :: OpsEnv -> Text -> StreamVersion -> Int -> IO OpsOutcome+runShow env name from limit =+ runAction env action $ \(info, events) -> Succeeded (streamResult name info events)+ where+ action = do+ info <- getStream (StreamName name)+ events <- readStreamForward (StreamName name) from (fromIntegral limit)+ pure (info, Vector.toList events)++runLifecycle :: OpsEnv -> Text -> Text -> Eff '[Store, Error StoreError, IOE] (Maybe StreamId) -> IO OpsOutcome+runLifecycle env operation name action+ | not env.force =+ runAction env (getStream (StreamName name)) $ \info ->+ PreviewRequired+ (lifecycleResult True operation name info Nothing)+ (forceInvocation env ["stream", operation, name])+ | otherwise =+ runAction env mutation $ \(before, outcome) ->+ Succeeded (lifecycleResult False operation name before outcome)+ where+ mutation = do+ before <- getStream (StreamName name)+ outcome <- action+ pure (before, outcome)++runHardDelete :: OpsEnv -> Text -> IO OpsOutcome+runHardDelete env name+ | not env.force = runLifecycle env "hard-delete" name (hardDeleteStream (StreamName name))+ | otherwise = do+ confirmed <- confirmDestructive env name+ if confirmed+ then runLifecycle env "hard-delete" name (hardDeleteStream (StreamName name))+ else pure (Failed "stream-name confirmation did not match; hard delete cancelled")++runTruncate :: OpsEnv -> TruncateCommand -> IO OpsOutcome+runTruncate env = \case+ ClearTruncateBefore name -> runClearTruncate env name+ SetTruncateBefore name before expected skipPreflight -> runSetTruncate env name before expected skipPreflight++runClearTruncate :: OpsEnv -> Text -> IO OpsOutcome+runClearTruncate env name+ | not env.force =+ runAction env (getStream (StreamName name)) $ \streamInfo ->+ PreviewRequired+ (lifecycleResult True "truncate-before clear" name streamInfo Nothing)+ (forceInvocation env ["stream", "truncate-before", "clear", name])+ | otherwise =+ runAction env action $ \(streamInfo, outcome) ->+ Succeeded (lifecycleResult False "truncate-before clear" name streamInfo outcome)+ where+ action = do+ streamInfo <- getStream (StreamName name)+ outcome <- clearStreamTruncateBefore (StreamName name)+ pure (streamInfo, outcome)++runSetTruncate :: OpsEnv -> Text -> StreamVersion -> Maybe Snapshot.ExpectedDiscriminators -> Bool -> IO OpsOutcome+runSetTruncate env name before expected skipPreflight = do+ checked <-+ if skipPreflight+ then pure Nothing+ else do+ outcome <- runStoreIO env.store (Snapshot.preflightFor name before expected)+ case outcome of+ Left storeError -> pure (Just (Left (Text.pack (show storeError))))+ Right evidence -> pure (Just (Right evidence))+ case checked of+ Just (Left message) -> pure (Failed message)+ Just (Right evidence)+ | not evidence.passed ->+ pure (Failed ("snapshot preflight failed: " <> evidence.reason))+ _+ | not env.force ->+ runAction env (getStream (StreamName name)) $ \info ->+ PreviewRequired+ (truncateResult True name before info checked)+ (forceInvocation env (setArguments name before expected skipPreflight))+ | otherwise -> do+ confirmed <- confirmDestructive env name+ if not confirmed+ then pure (Failed "stream-name confirmation did not match; truncate-before cancelled")+ else runAction env action $ \(info, outcome) ->+ Succeeded (truncateMutationResult name before info outcome checked)+ where+ action = do+ info <- getStream (StreamName name)+ outcome <- setStreamTruncateBefore (StreamName name) before+ pure (info, outcome)++setArguments :: Text -> StreamVersion -> Maybe Snapshot.ExpectedDiscriminators -> Bool -> [Text]+setArguments name before expected skipPreflight =+ ["stream", "truncate-before", "set", name, versionText before]+ <> maybe [] discriminatorArguments expected+ <> ["--skip-preflight" | skipPreflight]++discriminatorArguments :: Snapshot.ExpectedDiscriminators -> [Text]+discriminatorArguments expected =+ [ "--state-codec-version",+ showText expected.stateCodecVersion,+ "--regfile-shape-hash",+ expected.regfileShapeHash,+ "--state-shape-hash",+ expected.stateShapeHash+ ]++confirmDestructive :: OpsEnv -> Text -> IO Bool+confirmDestructive env name+ | env.outputMode == Json = pure True+ | otherwise = do+ Text.IO.putStr ("type the stream name to confirm: " <> name <> "\n> ")+ hFlush stdout+ entered <- Text.IO.getLine+ pure (entered == name)++runCausation :: OpsEnv -> EventId -> IO OpsOutcome+runCausation env eventId =+ runAction env action $ \(ancestors, descendants, names) ->+ Succeeded (causationResult eventId ancestors descendants names)+ where+ action = do+ ancestors <- Vector.toList <$> findCausationAncestors eventId+ descendants <- Vector.toList <$> findCausationDescendants eventId+ names <- lookupStreamNames (map (.originalStreamId) (ancestors <> descendants))+ pure (ancestors, descendants, names)++runAction :: OpsEnv -> Eff '[Store, Error StoreError, IOE] a -> (a -> OpsOutcome) -> IO OpsOutcome+runAction env action onSuccess = do+ result <- runStoreIO env.store action+ pure $ either (Failed . Text.pack . show) onSuccess result++streamResult :: Text -> Maybe StreamInfo -> [RecordedEvent] -> OpsResult+streamResult name info events =+ OpsResult+ { headers = ["stream", "version", "event_id", "event_type", "global_position", "created_at", "payload"],+ rows = map eventRow events,+ jsonValue = object ["stream" .= streamInfoJson name info, "events" .= map (eventJson Nothing) events]+ }+ where+ eventRow event =+ [ name,+ versionText event.streamVersion,+ eventIdText event.eventId,+ eventTypeText event.eventType,+ positionText event.globalPosition,+ showText event.createdAt,+ truncateCell 120 (jsonText event.payload)+ ]++streamInfoJson :: Text -> Maybe StreamInfo -> Value+streamInfoJson name = \case+ Nothing -> object ["name" .= name, "exists" .= False]+ Just info ->+ object+ [ "name" .= name,+ "exists" .= True,+ "stream_id" .= streamIdInt info.id,+ "version" .= versionInt info.version,+ "created_at" .= info.createdAt,+ "deleted_at" .= info.deletedAt,+ "truncate_before" .= versionInt info.truncateBefore+ ]++eventJson :: Maybe Text -> RecordedEvent -> Value+eventJson direction event =+ object+ [ "direction" .= direction,+ "event_id" .= eventIdText event.eventId,+ "event_type" .= eventTypeText event.eventType,+ "stream_version" .= versionInt event.streamVersion,+ "global_position" .= positionInt event.globalPosition,+ "original_stream_id" .= streamIdInt event.originalStreamId,+ "original_version" .= versionInt event.originalVersion,+ "payload" .= event.payload,+ "metadata" .= event.metadata,+ "causation_id" .= fmap UUID.toText event.causationId,+ "correlation_id" .= fmap UUID.toText event.correlationId,+ "created_at" .= event.createdAt+ ]++lifecycleResult :: Bool -> Text -> Text -> Maybe StreamInfo -> Maybe StreamId -> OpsResult+lifecycleResult preview operation name info outcome =+ OpsResult+ { headers = ["operation", "stream", "current_state", "disposition"],+ rows = [[operation, name, maybe "not_found" streamState info, disposition]],+ jsonValue = object ["preview" .= preview, "operation" .= operation, "stream" .= streamInfoJson name info, "affected_stream_id" .= fmap streamIdInt outcome, "disposition" .= disposition]+ }+ where+ disposition+ | preview = maybe "not_found" (const ("would_" <> Text.replace " " "_" operation)) info+ | otherwise = maybe "not_transitioned" (const "transitioned") outcome++streamState :: StreamInfo -> Text+streamState info+ | info.deletedAt == Nothing = "live"+ | otherwise = "soft_deleted"++truncateResult :: Bool -> Text -> StreamVersion -> Maybe StreamInfo -> Maybe (Either Text Snapshot.PreflightEvidence) -> OpsResult+truncateResult preview name before info checked =+ OpsResult+ { headers = ["stream", "current_marker", "new_marker", "preflight", "disposition"],+ rows = [[name, maybe "not_found" (versionText . (.truncateBefore)) info, versionText before, preflightText checked, if preview then "would_set" else "set"]],+ jsonValue = object ["preview" .= preview, "stream" .= streamInfoJson name info, "new_marker" .= versionInt before, "preflight" .= preflightJson checked]+ }++truncateMutationResult :: Text -> StreamVersion -> Maybe StreamInfo -> Maybe StreamId -> Maybe (Either Text Snapshot.PreflightEvidence) -> OpsResult+truncateMutationResult name before info outcome checked =+ (truncateResult False name before info checked)+ { rows = [[name, maybe "not_found" (versionText . (.truncateBefore)) info, versionText before, preflightText checked, maybe "not_transitioned" (const "set") outcome]],+ jsonValue = object ["preview" .= False, "stream" .= streamInfoJson name info, "new_marker" .= versionInt before, "preflight" .= preflightJson checked, "affected_stream_id" .= fmap streamIdInt outcome]+ }++preflightText :: Maybe (Either Text Snapshot.PreflightEvidence) -> Text+preflightText Nothing = "skipped"+preflightText (Just (Left message)) = "error: " <> message+preflightText (Just (Right evidence)) = if evidence.passed then "passed" else "failed: " <> evidence.reason++preflightJson :: Maybe (Either Text Snapshot.PreflightEvidence) -> Value+preflightJson Nothing = object ["skipped" .= True]+preflightJson (Just (Left message)) = object ["error" .= message]+preflightJson (Just (Right evidence)) =+ object+ [ "passed" .= evidence.passed,+ "reason" .= evidence.reason,+ "required_snapshot_version" .= versionInt evidence.requiredSnapshotVersion,+ "version_covered" .= evidence.versionCovered,+ "discriminators_match" .= evidence.discriminatorsMatch+ ]++causationResult :: EventId -> [RecordedEvent] -> [RecordedEvent] -> Map StreamId StreamName -> OpsResult+causationResult seed ancestors descendants names =+ OpsResult+ { headers = ["direction", "event_id", "stream", "type", "global_position", "causation_id"],+ rows = map (uncurry row) directed,+ jsonValue = Aeson.toJSON [eventJson (Just direction) event | (direction, event) <- directed]+ }+ where+ seedRows = take 1 (filter ((== seed) . (.eventId)) (ancestors <> descendants))+ directed =+ map ("seed",) seedRows+ <> map ("ancestor",) (filter ((/= seed) . (.eventId)) ancestors)+ <> map ("descendant",) (filter ((/= seed) . (.eventId)) descendants)+ row direction event =+ [ direction,+ eventIdText event.eventId,+ maybe ("#" <> showText (streamIdInt event.originalStreamId)) streamNameText (Map.lookup event.originalStreamId names),+ eventTypeText event.eventType,+ positionText event.globalPosition,+ maybe "" UUID.toText event.causationId+ ]++streamNameText :: StreamName -> Text+streamNameText (StreamName value) = value++eventIdText :: EventId -> Text+eventIdText (EventId value) = UUID.toText value++eventTypeText :: EventType -> Text+eventTypeText (EventType value) = value++streamIdInt :: StreamId -> Int64+streamIdInt (StreamId value) = value++versionInt :: StreamVersion -> Int64+versionInt (StreamVersion value) = value++positionInt :: GlobalPosition -> Int64+positionInt (GlobalPosition value) = value++versionText :: StreamVersion -> Text+versionText = showText . versionInt++positionText :: GlobalPosition -> Text+positionText = showText . positionInt++showText :: (Show a) => a -> Text+showText = Text.pack . show++forceInvocation :: OpsEnv -> [Text] -> Text+forceInvocation env arguments = Text.unwords (map shellQuote ("keiro-ops" : arguments <> globalFlags <> ["--force"]))+ where+ globalFlags = ["--json" | env.outputMode == Json] <> ["--allow-schema-drift" | env.allowSchemaDrift]++shellQuote :: Text -> Text+shellQuote value = "'" <> Text.replace "'" "'\"'\"'" value <> "'"
+ src/Keiro/Ops/Timer.hs view
@@ -0,0 +1,339 @@+module Keiro.Ops.Timer+ ( Command (..),+ DrainOptions (..),+ StuckListOptions (..),+ TimerFire,+ commandParser,+ commandParserWithDrain,+ isMutation,+ runCommand,+ runCommandWithFire,+ )+where++import Data.Aeson (Value, object, (.=))+import Data.Aeson qualified as Aeson+import Data.Text (Text)+import Data.Text qualified as Text+import Data.Time (NominalDiffTime, UTCTime, getCurrentTime)+import Data.UUID (UUID)+import Data.UUID qualified as UUID+import Effectful (Eff, IOE)+import Effectful.Error.Static (Error)+import Keiro.Ops.Env (OpsEnv (..), OutputMode (..))+import Keiro.Ops.Parse (durationReader, nonNegativeIntReader, positiveIntReader)+import Keiro.Ops.Render+import Keiro.Timer+import Kiroku.Store.Effect (Store, runStoreIO)+import Kiroku.Store.Error (StoreError)+import Kiroku.Store.Types (EventId (..))+import Options.Applicative hiding (action, value)+import Options.Applicative qualified as Optparse++data StuckListOptions = StuckListOptions+ { minAge :: !(Maybe NominalDiffTime),+ minAttempts :: !(Maybe Int)+ }+ deriving stock (Eq, Show)++data DrainOptions = DrainOptions+ { limit :: !Int+ }+ deriving stock (Eq, Show)++type TimerFire = TimerRow -> Eff '[Store, Error StoreError, IOE] (Maybe EventId)++data Command+ = StuckList !StuckListOptions+ | Requeue !TimerId+ | Cancel !TimerId+ | DeadLetter !TimerId !Text+ | DrainOnce !DrainOptions+ deriving stock (Eq, Show)++commandParser :: Parser Command+commandParser = commandParserWithDrain False++commandParserWithDrain :: Bool -> Parser Command+commandParserWithDrain includeDrain =+ hsubparser+ ( command+ "stuck"+ ( info+ stuckCommandParser+ (progDesc "Find timers left firing by an interrupted worker")+ )+ <> command+ "requeue"+ ( info+ (Requeue <$> timerIdArgument)+ (progDesc "Preview or return a stuck firing timer to the scheduled queue")+ )+ <> command+ "cancel"+ ( info+ (Cancel <$> timerIdArgument)+ (progDesc "Preview or withdraw a scheduled or firing timer permanently")+ )+ <> command+ "dead-letter"+ ( info+ ( DeadLetter+ <$> timerIdArgument+ <*> (Text.pack <$> strOption (long "reason" <> metavar "TEXT" <> help "Operator reason recorded with the terminal transition"))+ )+ (progDesc "Preview or abandon a scheduled or firing timer with a reason")+ )+ <> drainCommand+ )+ where+ drainCommand+ | includeDrain =+ command+ "drain-once"+ ( info+ (DrainOnce . DrainOptions <$> option positiveIntReader (long "limit" <> metavar "N" <> Optparse.value 100 <> showDefault <> help "Maximum due timers to dispatch"))+ (progDesc "Preview or run one bounded pass through the application timer-fire hook")+ )+ | otherwise = mempty++stuckCommandParser :: Parser Command+stuckCommandParser =+ hsubparser+ ( command+ "list"+ ( info+ ( StuckList+ <$> ( StuckListOptions+ <$> optional (option durationReader (long "min-age" <> metavar "DURATION" <> help "Minimum time in firing, such as 5m or 1h"))+ <*> optional (option nonNegativeIntReader (long "min-attempts" <> metavar "N" <> help "Minimum claim attempt count"))+ )+ )+ (progDesc "List stuck timers; requeue transient failures, cancel obsolete work, or dead-letter poison work")+ )+ )++timerIdArgument :: Parser TimerId+timerIdArgument = TimerId <$> argument uuidReader (metavar "TIMER_ID")++uuidReader :: ReadM UUID+uuidReader = eitherReader $ \raw ->+ maybe (Left "expected a UUID timer id") Right (UUID.fromString raw)++isMutation :: Command -> Bool+isMutation = \case+ StuckList {} -> False+ Requeue {} -> True+ Cancel {} -> True+ DeadLetter {} -> True+ DrainOnce {} -> True++runCommand :: OpsEnv -> Command -> IO OpsOutcome+runCommand = runCommandWithFire Nothing++runCommandWithFire :: Maybe TimerFire -> OpsEnv -> Command -> IO OpsOutcome+runCommandWithFire timerFire env = \case+ StuckList options -> runStuckList env options+ Requeue timerId -> runMutation env RequeueOperation timerId Nothing+ Cancel timerId -> runMutation env CancelOperation timerId Nothing+ DeadLetter timerId reason -> runMutation env DeadLetterOperation timerId (Just reason)+ DrainOnce options -> runDrainOnce timerFire env options++runDrainOnce :: Maybe TimerFire -> OpsEnv -> DrainOptions -> IO OpsOutcome+runDrainOnce Nothing _ _ = pure (Failed "timer fire hook is not mounted")+runDrainOnce (Just fire) env options = do+ now <- getCurrentTime+ if env.force+ then+ runAction+ env+ (drainDueTimersWith Nothing defaultTimerWorkerOptions now options.limit fire)+ (Succeeded . drainResult options.limit)+ else runAction env (countDueTimers now) $ \due ->+ PreviewRequired+ (drainPreviewResult options.limit due)+ (forceInvocation env ["timer", "drain-once", "--limit", Text.pack (show options.limit)])++drainPreviewResult :: Int -> Int -> OpsResult+drainPreviewResult limit due =+ OpsResult+ { headers = ["due", "limit", "would_process"],+ rows = [[showText due, showText limit, showText (min (toInteger due) (toInteger limit))]],+ jsonValue =+ object+ [ "preview" .= True,+ "due" .= due,+ "limit" .= limit,+ "would_process" .= min (toInteger due) (toInteger limit)+ ]+ }++drainResult :: Int -> Int -> OpsResult+drainResult limit processed =+ OpsResult+ { headers = ["limit", "processed"],+ rows = [[showText limit, showText processed]],+ jsonValue = object ["limit" .= limit, "processed" .= processed]+ }++showText :: (Show a) => a -> Text+showText = Text.pack . show++runStuckList :: OpsEnv -> StuckListOptions -> IO OpsOutcome+runStuckList env options = do+ now <- getCurrentTime+ runAction env (findStuckTimers now stuckFilter) (Succeeded . timerListResult)+ where+ stuckFilter = StuckTimerFilter options.minAge options.minAttempts++data TimerOperation+ = RequeueOperation+ | CancelOperation+ | DeadLetterOperation++runMutation :: OpsEnv -> TimerOperation -> TimerId -> Maybe Text -> IO OpsOutcome+runMutation env operation timerId reason+ | not env.force =+ runAction env (lookupTimer timerId) $ \row ->+ PreviewRequired+ (timerPreviewResult operation timerId row)+ (forceInvocation env (operationArguments operation timerId reason))+ | otherwise =+ runAction env action $ \(transitioned, row) ->+ Succeeded (timerMutationResult operation timerId transitioned row)+ where+ action = do+ transitioned <- case operation of+ RequeueOperation -> requeueStuckTimer timerId+ CancelOperation -> cancelTimer timerId+ DeadLetterOperation -> deadLetterTimer timerId (maybe "operator dead-letter" id reason)+ row <- lookupTimer timerId+ pure (transitioned, row)++runAction ::+ OpsEnv ->+ Eff '[Store, Error StoreError, IOE] a ->+ (a -> OpsOutcome) ->+ IO OpsOutcome+runAction env action onSuccess = do+ result <- runStoreIO env.store action+ pure $ case result of+ Left storeError -> Failed (Text.pack (show storeError))+ Right value -> onSuccess value++timerListResult :: [TimerRow] -> OpsResult+timerListResult timers =+ OpsResult+ { headers = ["id", "manager", "correlation", "fire_at", "status", "attempts", "payload"],+ rows = map timerRow timers,+ jsonValue = Aeson.toJSON (map timerJson timers)+ }++timerRow :: TimerRow -> [Text]+timerRow row =+ [ timerIdText row.timerId,+ row.processManagerName,+ row.correlationId,+ timeText row.fireAt,+ timerStatusText row.status,+ Text.pack (show row.attempts),+ truncateCell 120 (jsonText row.payload)+ ]++timerJson :: TimerRow -> Value+timerJson row =+ object+ [ "timer_id" .= timerIdText row.timerId,+ "process_manager_name" .= row.processManagerName,+ "correlation_id" .= row.correlationId,+ "fire_at" .= row.fireAt,+ "payload" .= row.payload,+ "status" .= timerStatusText row.status,+ "attempts" .= row.attempts,+ "fired_event_id" .= fmap (\(EventId eventId) -> UUID.toText eventId) row.firedEventId+ ]++timerPreviewResult :: TimerOperation -> TimerId -> Maybe TimerRow -> OpsResult+timerPreviewResult operation timerId row =+ OpsResult+ { headers = ["operation", "id", "disposition", "status"],+ rows = [[operationText operation, timerIdText timerId, disposition, maybe "not_found" (timerStatusText . (.status)) row]],+ jsonValue =+ object+ [ "preview" .= True,+ "operation" .= operationText operation,+ "disposition" .= disposition,+ "timer" .= fmap timerJson row+ ]+ }+ where+ disposition = timerDisposition operation row++timerMutationResult :: TimerOperation -> TimerId -> Bool -> Maybe TimerRow -> OpsResult+timerMutationResult operation timerId transitioned row =+ OpsResult+ { headers = ["operation", "id", "outcome", "status"],+ rows = [[operationText operation, timerIdText timerId, outcome, maybe "not_found" (timerStatusText . (.status)) row]],+ jsonValue =+ object+ [ "operation" .= operationText operation,+ "outcome" .= outcome,+ "transitioned" .= transitioned,+ "timer" .= fmap timerJson row+ ]+ }+ where+ outcome+ | transitioned = "transitioned"+ | otherwise = "not_transitioned"++timerDisposition :: TimerOperation -> Maybe TimerRow -> Text+timerDisposition _ Nothing = "not_found"+timerDisposition operation (Just row) = case operation of+ RequeueOperation+ | row.status == Firing -> "would_requeue"+ | otherwise -> "not_firing"+ CancelOperation+ | row.status `elem` [Scheduled, Firing] -> "would_cancel"+ | otherwise -> "already_terminal"+ DeadLetterOperation+ | row.status `elem` [Scheduled, Firing] -> "would_dead_letter"+ | otherwise -> "already_terminal"++operationArguments :: TimerOperation -> TimerId -> Maybe Text -> [Text]+operationArguments operation timerId reason =+ ["timer", operationText operation, timerIdText timerId]+ <> case operation of+ DeadLetterOperation -> ["--reason", maybe "operator dead-letter" id reason]+ _ -> []++operationText :: TimerOperation -> Text+operationText = \case+ RequeueOperation -> "requeue"+ CancelOperation -> "cancel"+ DeadLetterOperation -> "dead-letter"++timerStatusText :: TimerStatus -> Text+timerStatusText = \case+ Scheduled -> "scheduled"+ Firing -> "firing"+ Fired -> "fired"+ Cancelled -> "cancelled"+ Dead -> "dead"++timerIdText :: TimerId -> Text+timerIdText (TimerId timerId) = UUID.toText timerId++timeText :: UTCTime -> Text+timeText = Text.pack . show++forceInvocation :: OpsEnv -> [Text] -> Text+forceInvocation env arguments =+ Text.unwords (map shellQuote ("keiro-ops" : arguments <> globalFlags <> ["--force"]))+ where+ globalFlags =+ ["--json" | env.outputMode == Json]+ <> ["--allow-schema-drift" | env.allowSchemaDrift]++shellQuote :: Text -> Text+shellQuote value = "'" <> Text.replace "'" "'\"'\"'" value <> "'"
+ src/Keiro/Ops/Workflow.hs view
@@ -0,0 +1,882 @@+module Keiro.Ops.Workflow+ ( AwakeableCommand (..),+ Command (..),+ GcOptions (..),+ InspectOptions (..),+ ListOptions (..),+ PayloadArg (..),+ ResumeHook,+ ResumeOptions (..),+ WorkflowRef (..),+ commandParser,+ commandParserWithResume,+ isMutation,+ runCommand,+ runCommandWithResume,+ )+where++import Data.Aeson (Value, object, (.=))+import Data.Aeson qualified as Aeson+import Data.Aeson.Types qualified as AesonTypes+import Data.Int (Int64)+import Data.List.NonEmpty qualified as NonEmpty+import Data.Map.Strict (Map)+import Data.Map.Strict qualified as Map+import Data.Maybe (catMaybes)+import Data.Set qualified as Set+import Data.Text (Text)+import Data.Text qualified as Text+import Data.Text.Encoding qualified as Text.Encoding+import Data.Time (NominalDiffTime, UTCTime, getCurrentTime)+import Data.UUID (UUID)+import Data.UUID qualified as UUID+import Data.Vector qualified as Vector+import Effectful (Eff, IOE, (:>))+import Effectful.Error.Static (Error)+import Keiro.Codec (decodeRecorded)+import Keiro.Ops.Env (OpsEnv (..), OutputMode (..))+import Keiro.Ops.Parse (durationReader, nonNegativeReader, positiveIntReader)+import Keiro.Ops.Render+import Keiro.Workflow.Awakeable (AwakeableId (..), cancelAwakeable, signalAwakeable)+import Keiro.Workflow.Awakeable.Schema qualified as Awakeable+import Keiro.Workflow.Child.Schema qualified as Child+import Keiro.Workflow.Gc qualified as Gc+import Keiro.Workflow.Instance qualified as Instance+import Keiro.Workflow.Resume+ ( ResumeSummary (..),+ WorkflowRegistry,+ WorkflowResumeOptions,+ resumeWorkflowsOnceUpTo,+ )+import Keiro.Workflow.Schema qualified as WorkflowSchema+import Keiro.Workflow.Types+ ( WorkflowId (..),+ WorkflowJournalEvent (..),+ WorkflowName (..),+ awakeableAllocStepPrefix,+ cancelledStepName,+ completedStepName,+ continuedAsNewStepName,+ failedStepName,+ workflowGenerationStreamName,+ workflowJournalCodec,+ )+import Kiroku.Store.Effect (Store, runStoreIO)+import Kiroku.Store.Error (StoreError)+import Kiroku.Store.Read qualified as StoreRead+import Kiroku.Store.Types+import Options.Applicative hiding (action, value)+import Options.Applicative qualified as Optparse++data WorkflowRef = WorkflowRef+ { workflowName :: !Text,+ workflowId :: !Text+ }+ deriving stock (Eq, Show)++data ListOptions = ListOptions+ { statusFilters :: ![Instance.WorkflowStatus],+ workflowNameFilter :: !(Maybe Text),+ afterKey :: !(Maybe (Text, Text)),+ limit :: !Int+ }+ deriving stock (Eq, Show)++data InspectOptions = InspectOptions+ { target :: !WorkflowRef,+ generation :: !(Maybe Int)+ }+ deriving stock (Eq, Show)++data PayloadArg = PayloadArg+ { rawPayload :: !Text,+ payload :: !Value+ }+ deriving stock (Eq, Show)++data GcOptions = GcOptions+ { retention :: !NominalDiffTime,+ batchSize :: !Int+ }+ deriving stock (Eq, Show)++data ResumeOptions = ResumeOptions+ { limit :: !Int+ }+ deriving stock (Eq, Show)++type ResumeHook =+ ( WorkflowRegistry '[Store, Error StoreError, IOE],+ WorkflowResumeOptions+ )++data AwakeableCommand+ = AwakeableShow !UUID+ | AwakeableSignal !UUID !PayloadArg+ | AwakeableCancel !UUID+ deriving stock (Eq, Show)++data Command+ = List !ListOptions+ | Show !WorkflowRef+ | Steps !InspectOptions+ | Journal !InspectOptions+ | Awakeable !AwakeableCommand+ | Cancel !WorkflowRef+ | Resurrect !WorkflowRef+ | ReleaseLease !WorkflowRef+ | GcRunOnce !GcOptions+ | ResumeOnce !ResumeOptions+ deriving stock (Eq, Show)++commandParser :: Parser Command+commandParser = commandParserWithResume False++commandParserWithResume :: Bool -> Parser Command+commandParserWithResume includeResume =+ hsubparser+ ( command "list" (info (List <$> listOptionsParser) (progDesc "List workflow instances with stable keyset paging"))+ <> command "show" (info (Show <$> workflowRefParser) (progDesc "Show an instance, its children, and its awakeables"))+ <> command "steps" (info (Steps <$> inspectOptionsParser) (progDesc "Show the derived step index for one generation"))+ <> command "journal" (info (Journal <$> inspectOptionsParser) (progDesc "Decode one workflow journal generation in order"))+ <> command "awakeable" (info (Awakeable <$> awakeableCommandParser) (progDesc "Inspect, signal, or cancel an awakeable"))+ <> command "cancel" (info (Cancel <$> workflowRefParser) (progDesc "Preview or cancel a workflow at its next durable boundary"))+ <> command "resurrect" (info (Resurrect <$> workflowRefParser) (progDesc "Preview or resurrect a terminally failed workflow"))+ <> command "lease" (info leaseCommandParser (progDesc "Operate workflow instance leases"))+ <> command "gc" (info gcCommandParser (progDesc "Preview or run one workflow garbage-collection pass"))+ <> resumeCommand+ )+ where+ resumeCommand+ | includeResume =+ command+ "resume-once"+ ( info+ (ResumeOnce . ResumeOptions <$> option positiveIntReader (long "limit" <> metavar "N" <> Optparse.value 100 <> showDefault <> help "Maximum workflow instances to advance"))+ (progDesc "Preview or run one bounded application-registry resume pass")+ )+ | otherwise = mempty++listOptionsParser :: Parser ListOptions+listOptionsParser =+ ListOptions+ <$> many+ ( option+ workflowStatusReader+ (long "status" <> metavar "STATUS" <> help "Exact status; repeat to match more than one")+ )+ <*> optional (Text.pack <$> strOption (long "name" <> metavar "NAME" <> help "Exact workflow definition name"))+ <*> optional+ ( (,)+ <$> (Text.pack <$> strOption (long "after" <> metavar "NAME" <> help "Keyset cursor workflow name; followed by ID"))+ <*> (Text.pack <$> argument str (metavar "ID"))+ )+ <*> option positiveIntReader (long "limit" <> metavar "N" <> Optparse.value 100 <> showDefault <> help "Maximum rows to return")++inspectOptionsParser :: Parser InspectOptions+inspectOptionsParser =+ InspectOptions+ <$> workflowRefParser+ <*> optional (option (nonNegativeReader "expected a non-negative generation") (long "generation" <> metavar "N" <> help "Journal generation; defaults to current"))++workflowRefParser :: Parser WorkflowRef+workflowRefParser =+ WorkflowRef+ <$> (Text.pack <$> argument str (metavar "NAME"))+ <*> (Text.pack <$> argument str (metavar "ID"))++awakeableCommandParser :: Parser AwakeableCommand+awakeableCommandParser =+ hsubparser+ ( command "show" (info (AwakeableShow <$> uuidArgument) (progDesc "Show one awakeable"))+ <> command+ "signal"+ ( info+ (AwakeableSignal <$> uuidArgument <*> option payloadReader (long "payload" <> metavar "JSON" <> help "JSON completion payload"))+ (progDesc "Preview or signal a pending awakeable")+ )+ <> command "cancel" (info (AwakeableCancel <$> uuidArgument) (progDesc "Preview or cancel a pending awakeable"))+ )++leaseCommandParser :: Parser Command+leaseCommandParser =+ hsubparser+ (command "release" (info (ReleaseLease <$> workflowRefParser) (progDesc "Preview or forcibly release an instance lease")))++gcCommandParser :: Parser Command+gcCommandParser =+ hsubparser+ ( command+ "run-once"+ ( info+ ( GcRunOnce+ <$> ( GcOptions+ <$> option durationReader (long "retention" <> metavar "DURATION" <> help "Minimum terminal age, such as 30d or 12h")+ <*> option positiveIntReader (long "batch" <> metavar "N" <> Optparse.value 100 <> showDefault <> help "Maximum workflows to collect")+ )+ )+ (progDesc "Preview or run one bounded garbage-collection pass")+ )+ )++uuidArgument :: Parser UUID+uuidArgument = argument uuidReader (metavar "UUID")++uuidReader :: ReadM UUID+uuidReader = eitherReader $ \raw ->+ maybe (Left "expected a UUID") Right (UUID.fromString raw)++payloadReader :: ReadM PayloadArg+payloadReader = eitherReader $ \raw ->+ let rawText = Text.pack raw+ in case Aeson.eitherDecodeStrict' (Text.Encoding.encodeUtf8 rawText) of+ Left err -> Left ("invalid JSON payload: " <> err)+ Right value -> Right (PayloadArg rawText value)++workflowStatusReader :: ReadM Instance.WorkflowStatus+workflowStatusReader = eitherReader $ \case+ "running" -> Right Instance.WfRunning+ "suspended" -> Right Instance.WfSuspended+ "completed" -> Right Instance.WfCompleted+ "cancelled" -> Right Instance.WfCancelled+ "failed" -> Right Instance.WfFailed+ _ -> Left "expected one of: running, suspended, completed, cancelled, failed"++isMutation :: Command -> Bool+isMutation = \case+ List {} -> False+ Show {} -> False+ Steps {} -> False+ Journal {} -> False+ Awakeable (AwakeableShow {}) -> False+ Awakeable (AwakeableSignal {}) -> True+ Awakeable (AwakeableCancel {}) -> True+ Cancel {} -> True+ Resurrect {} -> True+ ReleaseLease {} -> True+ GcRunOnce {} -> True+ ResumeOnce {} -> True++runCommand :: OpsEnv -> Command -> IO OpsOutcome+runCommand = runCommandWithResume Nothing++runCommandWithResume :: Maybe ResumeHook -> OpsEnv -> Command -> IO OpsOutcome+runCommandWithResume resumeHook env = \case+ List options -> runList env options+ Show ref -> runShow env ref+ Steps options -> runSteps env options+ Journal options -> runJournal env options+ Awakeable awakeableCommand -> runAwakeable env awakeableCommand+ Cancel ref -> runCancel env ref+ Resurrect ref -> runResurrect env ref+ ReleaseLease ref -> runReleaseLease env ref+ GcRunOnce options -> runGc env options+ ResumeOnce options -> runResumeOnce resumeHook env options++runResumeOnce :: Maybe ResumeHook -> OpsEnv -> ResumeOptions -> IO OpsOutcome+runResumeOnce Nothing _ _ =+ pure (Failed "workflow resume hook is not mounted")+runResumeOnce (Just (registry, resumeOptions)) env options+ | not env.force = do+ now <- getCurrentTime+ runAction env (take options.limit <$> WorkflowSchema.findUnfinishedWorkflowIds now) $ \candidates ->+ PreviewRequired+ (resumePreviewResult options.limit candidates)+ (forceInvocation env ["wf", "resume-once", "--limit", Text.pack (show options.limit)])+ | otherwise =+ runAction+ env+ (resumeWorkflowsOnceUpTo options.limit resumeOptions registry)+ (Succeeded . resumeSummaryResult)++resumePreviewResult :: Int -> [(Text, Text)] -> OpsResult+resumePreviewResult limit candidates =+ OpsResult+ { headers = ["name", "id"],+ rows = [[name, workflowId] | (workflowId, name) <- candidates],+ jsonValue =+ object+ [ "preview" .= True,+ "limit" .= limit,+ "candidates"+ .= [ object ["workflow_name" .= name, "workflow_id" .= workflowId]+ | (workflowId, name) <- candidates+ ]+ ]+ }++resumeSummaryResult :: ResumeSummary -> OpsResult+resumeSummaryResult summary =+ OpsResult+ { headers = ["discovered", "resumed", "completed", "suspended", "unknown", "failed", "errors", "lease_skipped", "advanced", "paced", "sleep_due", "unregistered"],+ rows =+ [ map (Text.pack . show) counts+ <> [renderUnregistered summary.unregisteredNames]+ ],+ jsonValue =+ object+ [ "discovered" .= summary.discovered,+ "resumed" .= summary.resumed,+ "completed" .= summary.completed,+ "still_suspended" .= summary.stillSuspended,+ "unknown_name" .= summary.unknownName,+ "failed" .= summary.failed,+ "transient_errors" .= summary.transientErrors,+ "lease_skipped" .= summary.leaseSkipped,+ "advanced" .= summary.advanced,+ "paced" .= summary.paced,+ "sleep_due" .= summary.sleepDue,+ "unregistered_names" .= Set.toAscList summary.unregisteredNames+ ]+ }+ where+ counts =+ [ summary.discovered,+ summary.resumed,+ summary.completed,+ summary.stillSuspended,+ summary.unknownName,+ summary.failed,+ summary.transientErrors,+ summary.leaseSkipped,+ summary.advanced,+ summary.paced,+ summary.sleepDue+ ]+ renderUnregistered names+ | Set.null names = "-"+ | otherwise = Text.intercalate "," (Set.toAscList names)++runList :: OpsEnv -> ListOptions -> IO OpsOutcome+runList env options =+ runAction env (Instance.listWorkflowInstances filters) (Succeeded . workflowListResult)+ where+ filters =+ Instance.WorkflowInstanceFilter+ (NonEmpty.nonEmpty options.statusFilters)+ options.workflowNameFilter+ options.afterKey+ options.limit++runShow :: OpsEnv -> WorkflowRef -> IO OpsOutcome+runShow env ref =+ runAction env action $ \case+ Nothing -> Failed (workflowLabel ref <> " was not found")+ Just details -> Succeeded (workflowDetailsResult details)+ where+ action = do+ instanceRow <- Instance.lookupInstance (refName ref) (refId ref)+ case instanceRow of+ Nothing -> pure Nothing+ Just row -> do+ children <- Child.lookupChildrenOfParent ref.workflowId ref.workflowName+ awakeables <- lookupWorkflowAwakeables ref+ pure (Just (row, children, awakeables))++runSteps :: OpsEnv -> InspectOptions -> IO OpsOutcome+runSteps env options =+ runAction env action (Succeeded . uncurry stepsResult)+ where+ action = do+ selectedGeneration <- resolveGeneration options+ steps <- WorkflowSchema.loadStepIndex (refName options.target) (refId options.target) selectedGeneration+ pure (selectedGeneration, steps)++runJournal :: OpsEnv -> InspectOptions -> IO OpsOutcome+runJournal env options =+ runAction env action $ \(selectedGeneration, recorded) ->+ case traverse decodeJournalView recorded of+ Left err -> Failed err+ Right views -> Succeeded (journalResult selectedGeneration views)+ where+ action = do+ selectedGeneration <- resolveGeneration options+ recorded <- readJournalEvents (workflowGenerationStreamName (refName options.target) (refId options.target) selectedGeneration)+ pure (selectedGeneration, recorded)++runAwakeable :: OpsEnv -> AwakeableCommand -> IO OpsOutcome+runAwakeable env = \case+ AwakeableShow awakeableId ->+ runAction env (Awakeable.lookupAwakeable awakeableId) $ \case+ Nothing -> Failed ("awakeable " <> UUID.toText awakeableId <> " was not found")+ Just row -> Succeeded (awakeableResult row)+ AwakeableSignal awakeableId payloadArg+ | not env.force ->+ runAction env (Awakeable.lookupAwakeable awakeableId) $ \row ->+ PreviewRequired+ (awakeablePreviewResult "signal" awakeableId row)+ (forceInvocation env ["wf", "awakeable", "signal", UUID.toText awakeableId, "--payload", payloadArg.rawPayload])+ | otherwise ->+ runAction env action $ \(transitioned, row) ->+ Succeeded (awakeableMutationResult "signal" awakeableId transitioned row)+ where+ action = do+ transitioned <- signalAwakeable (AwakeableId awakeableId) payloadArg.payload+ row <- Awakeable.lookupAwakeable awakeableId+ pure (transitioned, row)+ AwakeableCancel awakeableId+ | not env.force ->+ runAction env (Awakeable.lookupAwakeable awakeableId) $ \row ->+ PreviewRequired+ (awakeablePreviewResult "cancel" awakeableId row)+ (forceInvocation env ["wf", "awakeable", "cancel", UUID.toText awakeableId])+ | otherwise ->+ runAction env action $ \(transitioned, row) ->+ Succeeded (awakeableMutationResult "cancel" awakeableId transitioned row)+ where+ action = do+ transitioned <- cancelAwakeable (AwakeableId awakeableId)+ row <- Awakeable.lookupAwakeable awakeableId+ pure (transitioned, row)++runCancel :: OpsEnv -> WorkflowRef -> IO OpsOutcome+runCancel env ref+ | not env.force =+ runAction env preview $ \(row, journalExists) ->+ PreviewRequired+ (instancePreviewResult "cancel" ref (cancelPreviewDisposition row journalExists) row)+ (forceInvocation env ["wf", "cancel", ref.workflowName, ref.workflowId])+ | otherwise =+ runAction env (Instance.cancelWorkflow (refName ref) (refId ref)) $ \outcome ->+ Succeeded (workflowMutationResult "cancel" ref (cancelOutcomeText outcome))+ where+ preview = do+ row <- Instance.lookupInstance (refName ref) (refId ref)+ journalExists <- case row of+ Just _ -> pure True+ Nothing -> do+ generation <- WorkflowSchema.currentGeneration (refName ref) (refId ref)+ not . Map.null <$> WorkflowSchema.loadStepIndex (refName ref) (refId ref) generation+ pure (row, journalExists)++runResurrect :: OpsEnv -> WorkflowRef -> IO OpsOutcome+runResurrect env ref+ | not env.force =+ runAction env (Instance.lookupInstance (refName ref) (refId ref)) $ \row ->+ PreviewRequired+ (instancePreviewResult "resurrect" ref (resurrectPreviewDisposition row) row)+ (forceInvocation env ["wf", "resurrect", ref.workflowName, ref.workflowId])+ | otherwise =+ runAction env (Instance.resurrectFailedWorkflow (refName ref) (refId ref)) $ \outcome ->+ Succeeded (workflowMutationResult "resurrect" ref (resurrectOutcomeText outcome))++runReleaseLease :: OpsEnv -> WorkflowRef -> IO OpsOutcome+runReleaseLease env ref+ | not env.force =+ runAction env (Instance.lookupInstance (refName ref) (refId ref)) $ \row ->+ PreviewRequired+ (instancePreviewResult "lease release" ref (leasePreviewDisposition row) row)+ (forceInvocation env ["wf", "lease", "release", ref.workflowName, ref.workflowId])+ | otherwise =+ runAction env (Instance.forceReleaseInstanceLease (refName ref) (refId ref)) $ \released ->+ Succeeded (workflowMutationResult "lease release" ref (if released then "released" else "no_lease_released"))++runGc :: OpsEnv -> GcOptions -> IO OpsOutcome+runGc env options+ | not env.force = do+ now <- getCurrentTime+ runAction env (Gc.listWorkflowGcCandidates now policy) $ \candidates ->+ PreviewRequired+ (gcCandidatesResult candidates)+ ( forceInvocation+ env+ [ "wf",+ "gc",+ "run-once",+ "--retention",+ Text.pack (show (realToFrac options.retention :: Double)) <> "s",+ "--batch",+ Text.pack (show options.batchSize)+ ]+ )+ | otherwise = do+ now <- getCurrentTime+ runAction env (Gc.gcWorkflowsOnce now policy) (Succeeded . gcSummaryResult)+ where+ policy = Gc.WorkflowGcPolicy options.retention options.batchSize++runAction ::+ OpsEnv ->+ Eff '[Store, Error StoreError, IOE] a ->+ (a -> OpsOutcome) ->+ IO OpsOutcome+runAction env action onSuccess = do+ result <- runStoreIO env.store action+ pure $ case result of+ Left storeError -> Failed (Text.pack (show storeError))+ Right value -> onSuccess value++resolveGeneration :: (Store :> es) => InspectOptions -> Eff es Int+resolveGeneration options =+ maybe+ (WorkflowSchema.currentGeneration (refName options.target) (refId options.target))+ pure+ options.generation++lookupWorkflowAwakeables :: (Store :> es) => WorkflowRef -> Eff es [Awakeable.AwakeableRow]+lookupWorkflowAwakeables ref = do+ current <- WorkflowSchema.currentGeneration (refName ref) (refId ref)+ stepIndexes <- traverse (WorkflowSchema.loadStepIndex (refName ref) (refId ref)) [0 .. current]+ catMaybes <$> traverse Awakeable.lookupAwakeable (awakeableIds stepIndexes)++awakeableIds :: [Map Text Value] -> [UUID]+awakeableIds indexes =+ Set.toAscList . Set.fromList $ do+ index <- indexes+ (stepName, value) <- Map.toList index+ if awakeableAllocStepPrefix `Text.isPrefixOf` stepName+ then case Aeson.fromJSON value of+ AesonTypes.Success (AwakeableId awakeableId) -> [awakeableId]+ AesonTypes.Error _ -> []+ else []++readJournalEvents :: (Store :> es) => StreamName -> Eff es [RecordedEvent]+readJournalEvents streamName = go (StreamVersion 0) []+ where+ pageSize = 256+ go cursor pages = do+ page <- StoreRead.readStreamForward streamName cursor pageSize+ if Vector.null page+ then pure (concat (reverse pages))+ else+ let nextCursor = (Vector.last page).streamVersion+ in go nextCursor (Vector.toList page : pages)++data JournalView = JournalView+ { eventId :: !Text,+ eventType :: !Text,+ streamVersion :: !Int64,+ globalPosition :: !Int64,+ stepName :: !Text,+ recordedAt :: !UTCTime,+ payload :: !Value+ }++decodeJournalView :: RecordedEvent -> Either Text JournalView+decodeJournalView recorded = do+ event <- firstShow (decodeRecorded workflowJournalCodec recorded)+ let (stepName, recordedAt, payload) = case event of+ StepRecorded name value timestamp -> (name, timestamp, value)+ WorkflowCompleted timestamp -> (completedStepName, timestamp, Aeson.Null)+ WorkflowCancelled timestamp -> (cancelledStepName, timestamp, Aeson.Null)+ WorkflowFailed reason timestamp -> (failedStepName, timestamp, Aeson.toJSON reason)+ WorkflowContinuedAsNew generation timestamp -> (continuedAsNewStepName, timestamp, Aeson.toJSON generation)+ pure+ JournalView+ { eventId = case recorded.eventId of EventId value -> UUID.toText value,+ eventType = case recorded.eventType of EventType value -> value,+ streamVersion = case recorded.streamVersion of StreamVersion value -> value,+ globalPosition = case recorded.globalPosition of GlobalPosition value -> value,+ stepName,+ recordedAt,+ payload+ }++firstShow :: (Show err) => Either err value -> Either Text value+firstShow = \case+ Left err -> Left ("workflow journal decode failed: " <> Text.pack (show err))+ Right value -> Right value++workflowListResult :: [Instance.WorkflowInstanceRow] -> OpsResult+workflowListResult instances =+ OpsResult+ { headers = ["name", "id", "generation", "status", "attempts", "lease", "wake_after", "updated_at"],+ rows = map workflowListRow instances,+ jsonValue = Aeson.toJSON (map workflowInstanceJson instances)+ }++workflowListRow :: Instance.WorkflowInstanceRow -> [Text]+workflowListRow row =+ [ row.workflowName,+ row.workflowId,+ Text.pack (show row.generation),+ Instance.statusToText row.status,+ Text.pack (show row.attempts),+ leaseText row,+ maybeTime row.wakeAfter,+ timeText row.updatedAt+ ]++workflowInstanceJson :: Instance.WorkflowInstanceRow -> Value+workflowInstanceJson row =+ object+ [ "workflow_id" .= row.workflowId,+ "workflow_name" .= row.workflowName,+ "generation" .= row.generation,+ "status" .= Instance.statusToText row.status,+ "attempts" .= row.attempts,+ "last_error" .= row.lastError,+ "next_attempt_at" .= row.nextAttemptAt,+ "wake_after" .= row.wakeAfter,+ "leased_by" .= row.leasedBy,+ "lease_expires_at" .= row.leaseExpiresAt,+ "created_at" .= row.createdAt,+ "updated_at" .= row.updatedAt,+ "completed_at" .= row.completedAt+ ]++workflowDetailsResult :: (Instance.WorkflowInstanceRow, [Child.ChildRow], [Awakeable.AwakeableRow]) -> OpsResult+workflowDetailsResult (row, children, awakeables) =+ OpsResult+ { headers = ["name", "id", "generation", "status", "attempts", "lease", "wake_after", "children", "awakeables"],+ rows =+ [ [ row.workflowName,+ row.workflowId,+ Text.pack (show row.generation),+ Instance.statusToText row.status,+ Text.pack (show row.attempts),+ leaseText row,+ maybeTime row.wakeAfter,+ Text.pack (show (length children)),+ Text.pack (show (length awakeables))+ ]+ ],+ jsonValue =+ object+ [ "instance" .= workflowInstanceJson row,+ "children" .= map childJson children,+ "awakeables" .= map awakeableJson awakeables+ ]+ }++childJson :: Child.ChildRow -> Value+childJson row =+ object+ [ "child_id" .= row.childId,+ "child_name" .= row.childName,+ "parent_id" .= row.parentId,+ "parent_name" .= row.parentName,+ "await_step" .= row.awaitStep,+ "status" .= Child.statusToText row.status,+ "result" .= row.result,+ "failure_reason" .= row.failureReason,+ "created_at" .= row.createdAt,+ "updated_at" .= row.updatedAt,+ "completed_at" .= row.completedAt+ ]++stepsResult :: Int -> Map Text Value -> OpsResult+stepsResult generation steps =+ OpsResult+ { headers = ["step", "result"],+ rows = [[name, truncateCell 120 (jsonText value)] | (name, value) <- Map.toAscList steps],+ jsonValue =+ object+ [ "generation" .= generation,+ "steps" .= [object ["step" .= name, "result" .= value] | (name, value) <- Map.toAscList steps]+ ]+ }++journalResult :: Int -> [JournalView] -> OpsResult+journalResult generation views =+ OpsResult+ { headers = ["version", "event_type", "step", "recorded_at", "payload"],+ rows =+ [ [ Text.pack (show view.streamVersion),+ view.eventType,+ view.stepName,+ timeText view.recordedAt,+ truncateCell 120 (jsonText view.payload)+ ]+ | view <- views+ ],+ jsonValue = object ["generation" .= generation, "events" .= map journalViewJson views]+ }++journalViewJson :: JournalView -> Value+journalViewJson view =+ object+ [ "event_id" .= view.eventId,+ "event_type" .= view.eventType,+ "stream_version" .= view.streamVersion,+ "global_position" .= view.globalPosition,+ "step_name" .= view.stepName,+ "recorded_at" .= view.recordedAt,+ "payload" .= view.payload+ ]++awakeableResult :: Awakeable.AwakeableRow -> OpsResult+awakeableResult row =+ OpsResult+ { headers = ["id", "owner_name", "owner_id", "status", "payload", "updated_at"],+ rows =+ [ [ UUID.toText row.awakeableId,+ row.ownerWorkflowName,+ row.ownerWorkflowId,+ Awakeable.statusToText row.status,+ maybe "-" (truncateCell 120 . jsonText) row.payload,+ timeText row.updatedAt+ ]+ ],+ jsonValue = awakeableJson row+ }++awakeableJson :: Awakeable.AwakeableRow -> Value+awakeableJson row =+ object+ [ "awakeable_id" .= UUID.toText row.awakeableId,+ "owner_workflow_name" .= row.ownerWorkflowName,+ "owner_workflow_id" .= row.ownerWorkflowId,+ "status" .= Awakeable.statusToText row.status,+ "payload" .= row.payload,+ "created_at" .= row.createdAt,+ "updated_at" .= row.updatedAt,+ "completed_at" .= row.completedAt+ ]++awakeablePreviewResult :: Text -> UUID -> Maybe Awakeable.AwakeableRow -> OpsResult+awakeablePreviewResult operation awakeableId row =+ OpsResult+ { headers = ["operation", "id", "disposition", "status"],+ rows = [[operation, UUID.toText awakeableId, disposition, maybe "not_found" (Awakeable.statusToText . (.status)) row]],+ jsonValue =+ object+ [ "preview" .= True,+ "operation" .= operation,+ "disposition" .= disposition,+ "awakeable" .= fmap awakeableJson row+ ]+ }+ where+ disposition = case row of+ Nothing -> "not_found"+ Just found -> case found.status of+ Awakeable.Pending -> "would_mutate"+ Awakeable.Completed | operation == "signal" -> "would_repair_if_needed"+ _ -> "no_op"++awakeableMutationResult :: Text -> UUID -> Bool -> Maybe Awakeable.AwakeableRow -> OpsResult+awakeableMutationResult operation awakeableId transitioned row =+ OpsResult+ { headers = ["operation", "id", "outcome", "status"],+ rows = [[operation, UUID.toText awakeableId, outcome, maybe "not_found" (Awakeable.statusToText . (.status)) row]],+ jsonValue =+ object+ [ "operation" .= operation,+ "outcome" .= outcome,+ "transitioned" .= transitioned,+ "awakeable" .= fmap awakeableJson row+ ]+ }+ where+ outcome+ | transitioned = "transitioned"+ | otherwise = "not_transitioned"++instancePreviewResult :: Text -> WorkflowRef -> Text -> Maybe Instance.WorkflowInstanceRow -> OpsResult+instancePreviewResult operation ref disposition row =+ OpsResult+ { headers = ["operation", "name", "id", "disposition", "status"],+ rows = [[operation, ref.workflowName, ref.workflowId, disposition, maybe "not_found" (Instance.statusToText . (.status)) row]],+ jsonValue =+ object+ [ "preview" .= True,+ "operation" .= operation,+ "disposition" .= disposition,+ "target" .= object ["workflow_name" .= ref.workflowName, "workflow_id" .= ref.workflowId],+ "instance" .= fmap workflowInstanceJson row+ ]+ }++workflowMutationResult :: Text -> WorkflowRef -> Text -> OpsResult+workflowMutationResult operation ref outcome =+ OpsResult+ { headers = ["operation", "name", "id", "outcome"],+ rows = [[operation, ref.workflowName, ref.workflowId, outcome]],+ jsonValue =+ object+ [ "operation" .= operation,+ "workflow_name" .= ref.workflowName,+ "workflow_id" .= ref.workflowId,+ "outcome" .= outcome+ ]+ }++gcCandidatesResult :: [Gc.WorkflowGcCandidate] -> OpsResult+gcCandidatesResult candidates =+ OpsResult+ { headers = ["name", "id", "disposition"],+ rows = [[candidate.workflowName, candidate.workflowId, "would_collect"] | candidate <- candidates],+ jsonValue =+ object+ [ "preview" .= True,+ "candidates"+ .= [ object ["workflow_name" .= candidate.workflowName, "workflow_id" .= candidate.workflowId]+ | candidate <- candidates+ ]+ ]+ }++gcSummaryResult :: Gc.WorkflowGcSummary -> OpsResult+gcSummaryResult summary =+ OpsResult+ { headers = ["scanned", "deleted"],+ rows = [[Text.pack (show summary.scanned), Text.pack (show summary.deleted)]],+ jsonValue = object ["scanned" .= summary.scanned, "deleted" .= summary.deleted]+ }++cancelPreviewDisposition :: Maybe Instance.WorkflowInstanceRow -> Bool -> Text+cancelPreviewDisposition row journalExists = case row of+ Just found -> case found.status of+ Instance.WfRunning -> "would_cancel"+ Instance.WfSuspended -> "would_cancel"+ _ -> "already_terminal"+ Nothing+ | journalExists -> "would_cancel_journal_only_instance"+ | otherwise -> "not_found"++resurrectPreviewDisposition :: Maybe Instance.WorkflowInstanceRow -> Text+resurrectPreviewDisposition = \case+ Just row | row.status == Instance.WfFailed -> "would_resurrect"+ Just _ -> "not_failed"+ Nothing -> "not_found"++leasePreviewDisposition :: Maybe Instance.WorkflowInstanceRow -> Text+leasePreviewDisposition = \case+ Just row | Just _ <- row.leasedBy -> "would_release"+ Just _ -> "no_lease"+ Nothing -> "not_found"++cancelOutcomeText :: Instance.CancelWorkflowOutcome -> Text+cancelOutcomeText = \case+ Instance.WorkflowCancelRecorded -> "cancel_recorded"+ Instance.WorkflowAlreadyTerminal status -> "already_" <> Instance.statusToText status+ Instance.WorkflowCancelUnknown -> "not_found"++resurrectOutcomeText :: Instance.ResurrectOutcome -> Text+resurrectOutcomeText = \case+ Instance.WorkflowResurrected -> "resurrected"+ Instance.WorkflowNotFailed -> "not_failed"+ Instance.WorkflowNotFound -> "not_found"++refName :: WorkflowRef -> WorkflowName+refName = WorkflowName . (.workflowName)++refId :: WorkflowRef -> WorkflowId+refId = WorkflowId . (.workflowId)++workflowLabel :: WorkflowRef -> Text+workflowLabel ref = ref.workflowName <> "/" <> ref.workflowId++leaseText :: Instance.WorkflowInstanceRow -> Text+leaseText row = case row.leasedBy of+ Nothing -> "-"+ Just owner -> owner <> maybe "" ((" until " <>) . timeText) row.leaseExpiresAt++maybeTime :: Maybe UTCTime -> Text+maybeTime = maybe "-" timeText++timeText :: UTCTime -> Text+timeText = Text.pack . show++forceInvocation :: OpsEnv -> [Text] -> Text+forceInvocation env arguments =+ Text.unwords (map shellQuote ("keiro-ops" : arguments <> globalFlags <> ["--force"]))+ where+ globalFlags =+ ["--json" | env.outputMode == Json]+ <> ["--allow-schema-drift" | env.allowSchemaDrift]++shellQuote :: Text -> Text+shellQuote value = "'" <> Text.replace "'" "'\"'\"'" value <> "'"
+ test/Main.hs view
@@ -0,0 +1,1729 @@+module Main (main) where++import Control.Exception (bracket)+import Data.Aeson (object, (.=))+import Data.Aeson qualified as Aeson+import Data.Aeson.Key (Key)+import Data.Aeson.KeyMap qualified as KeyMap+import Data.ByteString.Char8 qualified as ByteString+import Data.Either (isRight)+import Data.Function qualified as Function+import Data.Functor ((<&>))+import Data.Int (Int64)+import Data.List.NonEmpty qualified as NonEmpty+import Data.Map.Strict qualified as Map+import Data.Maybe (isJust)+import Data.Text (Text)+import Data.Text qualified as Text+import Data.Text.Encoding qualified as Text.Encoding+import Data.Time (UTCTime, addUTCTime, getCurrentTime)+import Data.UUID qualified as UUID+import Data.Vector qualified as Vector+import Effectful (Eff, IOE)+import Effectful.Error.Static (Error)+import Hasql.Connection qualified as Hasql+import Hasql.Connection.Settings qualified as HasqlSettings+import Hasql.Session qualified as HasqlSession+import Hasql.Transaction qualified as Tx+import Keiro.DeadLetter+import Keiro.Inbox qualified as Inbox+import Keiro.Integration.Event+import Keiro.Ops (AppHooks (..))+import Keiro.Ops qualified as Ops+import Keiro.Ops.Env+import Keiro.Ops.Inbox qualified as OpsInbox+import Keiro.Ops.Outbox qualified as OpsOutbox+import Keiro.Ops.Parse (parseDuration)+import Keiro.Ops.Pgmq qualified as OpsPgmq+import Keiro.Ops.Projection qualified as OpsProjection+import Keiro.Ops.Rebuild qualified as OpsRebuild+import Keiro.Ops.Render+import Keiro.Ops.ReplayAudit qualified as OpsReplayAudit+import Keiro.Ops.Shard qualified as OpsShard+import Keiro.Ops.Snapshot qualified as OpsSnapshot+import Keiro.Ops.Stream qualified as OpsStream+import Keiro.Ops.Timer qualified as OpsTimer+import Keiro.Ops.Workflow qualified as OpsWorkflow+import Keiro.Outbox qualified as Outbox+import Keiro.PGMQ+import Keiro.Projection qualified as Projection+import Keiro.Projection.Catalog qualified as Catalog+import Keiro.Projection.Catalog.Operations qualified as CatalogOperations+import Keiro.ReadModel.Rebuild qualified as Rebuild+import Keiro.Snapshot.Schema+import Keiro.Subscription.Shard qualified as Shard+import Keiro.Test.Postgres (Fixture, withFreshDatabase, withFreshStore, withMigratedSuiteWith)+import Keiro.Timer qualified as Timer+import Keiro.Workflow (StepName (..), WorkflowId (..), WorkflowJournalEvent (..), WorkflowName (..), appendJournalEntry)+import Keiro.Workflow.Awakeable (AwakeableId (..))+import Keiro.Workflow.Awakeable.Schema qualified as Awakeable+import Keiro.Workflow.Instance qualified as Instance+import Keiro.Workflow.Resume (WorkflowDef (..), defaultWorkflowResumeOptions)+import Keiro.Workflow.Sleep (sleepNamed)+import Kiroku.Store.Append (appendToStream)+import Kiroku.Store.Connection (KirokuStore (..))+import Kiroku.Store.Effect (Store, runStoreIO)+import Kiroku.Store.Error (StoreError)+import Kiroku.Store.HistoryRetention (StreamHistoryUnavailable (..))+import Kiroku.Store.Lifecycle (hardDeleteStream)+import Kiroku.Store.Read (getStream, readStreamForward)+import Kiroku.Store.Subscription.Types (SubscriptionName (..))+import Kiroku.Store.Transaction (runTransaction)+import Kiroku.Store.Types+import Options.Applicative qualified as Optparse+import Pgmq.Migration qualified as PgmqMigration+import System.Exit (ExitCode (..))+import System.Process (readProcessWithExitCode)+import Test.Hspec++main :: IO ()+main = do+ pgmq <- either (fail . show) pure PgmqMigration.pgmqMigrations+ withMigratedSuiteWith [pgmq] $ \fixture -> hspec (spec fixture)++embeddedHooks :: AppHooks+embeddedHooks =+ AppHooks+ { workflowResume = Just (Map.empty, defaultWorkflowResumeOptions),+ timerFire = Just (\_ -> pure Nothing),+ replayAudit = Just (OpsReplayAudit.OpsAuditConfig []),+ projectionCatalog = Just emptyCatalogOperations+ }++emptyCatalogOperations :: CatalogOperations.ProjectionCatalogOperations+emptyCatalogOperations =+ case Catalog.validateProjectionCatalog Catalog.emptyProjectionCatalog of+ Catalog.Failure diagnostics -> error ("empty projection catalog was invalid: " <> show diagnostics)+ Catalog.Success catalog -> CatalogOperations.projectionCatalogOperations catalog++parseOps :: AppHooks -> [String] -> Optparse.ParserResult Ops.OpsInvocation+parseOps hooks = Optparse.execParserPure Optparse.defaultPrefs (Ops.opsCommandTree hooks)++isParseSuccess :: Optparse.ParserResult value -> Bool+isParseSuccess Optparse.Success {} = True+isParseSuccess _ = False++isParseFailure :: Optparse.ParserResult value -> Bool+isParseFailure Optparse.Failure {} = True+isParseFailure _ = False++spec :: Fixture -> Spec+spec fixture = do+ describe "embedded command tree" do+ it "omits code-dependent commands from the standalone tree" do+ isParseFailure (parseOps Ops.emptyAppHooks ["wf", "resume-once"]) `shouldBe` True+ isParseFailure (parseOps Ops.emptyAppHooks ["timer", "drain-once"]) `shouldBe` True+ isParseFailure (parseOps Ops.emptyAppHooks ["replay-audit", "--full"]) `shouldBe` True+ isParseFailure (parseOps Ops.emptyAppHooks ["rebuild", "list"]) `shouldBe` True+ isParseFailure (parseOps Ops.emptyAppHooks ["rebuild", "adopt", "ops-group"]) `shouldBe` True+ isParseFailure (parseOps Ops.emptyAppHooks ["rebuild", "versioned", "status", "ops-run"]) `shouldBe` True+ isParseFailure (parseOps Ops.emptyAppHooks ["rebuild", "retired"]) `shouldBe` True+ isParseFailure (parseOps Ops.emptyAppHooks ["rebuild", "external-read", "counter_reader", "1"]) `shouldBe` True+ isParseFailure (parseOps Ops.emptyAppHooks ["rebuild", "reproject-stream", "ops-group", "ops-projection", "orders-1"]) `shouldBe` True++ it "mounts every code-dependent command from typed application hooks" do+ isParseSuccess (parseOps embeddedHooks ["wf", "resume-once"]) `shouldBe` True+ isParseSuccess (parseOps embeddedHooks ["timer", "drain-once"]) `shouldBe` True+ isParseSuccess (parseOps embeddedHooks ["replay-audit", "--full"]) `shouldBe` True+ isParseSuccess (parseOps embeddedHooks ["rebuild", "list"]) `shouldBe` True+ isParseSuccess (parseOps embeddedHooks ["rebuild", "adopt", "ops-group"]) `shouldBe` True+ isParseSuccess (parseOps embeddedHooks ["rebuild", "versioned", "status", "ops-run"]) `shouldBe` True+ isParseSuccess (parseOps embeddedHooks ["rebuild", "versioned", "resume", "ops-run"]) `shouldBe` True+ isParseSuccess (parseOps embeddedHooks ["rebuild", "versioned", "abandon", "ops-run"]) `shouldBe` True+ isParseSuccess (parseOps embeddedHooks ["rebuild", "retired"]) `shouldBe` True+ isParseSuccess (parseOps embeddedHooks ["rebuild", "drop-retired", "65b86cd6-550c-47c3-ae99-4039a85a11ad"]) `shouldBe` True+ isParseSuccess (parseOps embeddedHooks ["rebuild", "external-read", "counter_reader", "1"]) `shouldBe` True+ isParseSuccess (parseOps embeddedHooks ["rebuild", "retire-external-read", "counter_reader", "1"]) `shouldBe` True+ isParseSuccess (parseOps embeddedHooks ["rebuild", "reproject-stream", "ops-group", "ops-projection", "orders-1"]) `shouldBe` True+ isParseFailure (parseOps embeddedHooks ["rebuild", "retire-external-read", "counter_reader", "0"]) `shouldBe` True+ isParseFailure (parseOps embeddedHooks ["rebuild", "reproject-stream", "ops-group", "ops-projection", "orders-1", "--page-size", "0"]) `shouldBe` True++ it "parses a complete versioned start and rejects malformed generation identities" do+ let versionedStart =+ [ "rebuild",+ "versioned",+ "start",+ "ops-group",+ "--run-id",+ "ops-versioned-run",+ "--serving-revision",+ "revision-v1",+ "--candidate-revision",+ "revision-v2",+ "--target-mode",+ "clone",+ "--requested-by",+ "operator",+ "--reason",+ "schema repair"+ ]+ isParseSuccess (parseOps embeddedHooks versionedStart) `shouldBe` True+ isParseFailure (parseOps embeddedHooks ["rebuild", "drop-retired", "not-a-uuid"]) `shouldBe` True++ describe "numeric option rejection" do+ it "rejects non-finite durations on every duration flag" do+ isParseFailure (parseOps embeddedHooks ["outbox", "gc-sent", "--older-than", "NaN"]) `shouldBe` True+ isParseFailure (parseOps embeddedHooks ["outbox", "requeue-stuck", "--older-than", "Infinity"]) `shouldBe` True+ isParseFailure (parseOps embeddedHooks ["inbox", "gc", "--older-than", "NaN"]) `shouldBe` True+ isParseFailure (parseOps embeddedHooks ["timer", "stuck", "list", "--min-age", "NaNd"]) `shouldBe` True+ isParseFailure (parseOps embeddedHooks ["wf", "gc", "run-once", "--retention", "NaN", "--batch", "100"]) `shouldBe` True++ it "rejects non-positive and wrapped integer options at parse time" do+ isParseFailure (parseOps embeddedHooks ["wf", "gc", "run-once", "--retention", "30d", "--batch", "0"]) `shouldBe` True+ isParseFailure (parseOps embeddedHooks ["wf", "gc", "run-once", "--retention", "30d", "--batch=-5"]) `shouldBe` True+ isParseFailure (parseOps embeddedHooks ["wf", "list", "--limit", "0"]) `shouldBe` True+ isParseFailure (parseOps embeddedHooks ["replay-audit", "--full", "--resume-from", "-1"]) `shouldBe` True+ isParseFailure (parseOps embeddedHooks ["outbox", "list", "--source", "s", "--limit", "18446744073709551716"]) `shouldBe` True+ isParseSuccess (parseOps embeddedHooks ["wf", "gc", "run-once", "--retention", "30d", "--batch", "100"]) `shouldBe` True+ isParseSuccess (parseOps embeddedHooks ["replay-audit", "--full", "--resume-from", "0"]) `shouldBe` True++ it "shares non-negative admission across global positions, stream versions, and generations" do+ let rebuild position = ["rebuild", "start", "ops-group", "--run-id", "ops-run", "--requested-by", "test", "--reason", "test", "--from", position]+ snapshot version = ["snapshot", "truncation-preflight", "--stream", "orders-1", "--before", version]+ stream version = ["stream", "show", "orders-1", "--from", version]+ workflow generation = ["wf", "steps", "orders", "1", "--generation", generation]+ mapM_ (\args -> isParseFailure (parseOps embeddedHooks args) `shouldBe` True) [rebuild "-1", snapshot "-1", stream "-1", workflow "-1"]+ mapM_ (\args -> isParseSuccess (parseOps embeddedHooks args) `shouldBe` True) [rebuild "0", snapshot "0", stream "0", workflow "0"]++ describe "targeted stream repair command" $+ around (withFreshStore fixture) $ do+ it "requires a positive event admission limit" $ \_ -> do+ let command limit =+ [ "rebuild",+ "reproject-stream",+ "ops-group",+ "ops-projection",+ "ops-1",+ "--max-events",+ limit+ ]+ isParseSuccess (parseOps embeddedHooks (command "100")) `shouldBe` True+ isParseFailure (parseOps embeddedHooks (command "0")) `shouldBe` True++ it "classifies the command as mutating and renders a stable typed refusal" $ \store -> do+ let command =+ OpsRebuild.ReprojectStream+ OpsRebuild.ReprojectStreamOptions+ { groupId = either (error . show) Function.id (Catalog.mkRebuildGroupId "ops-group"),+ projectionId = either (error . show) Function.id (Catalog.mkProjectionId "ops-projection"),+ streamName = StreamName "orders-1",+ pageSize = 500,+ maxEvents = 1000+ }+ OpsRebuild.isMutation command `shouldBe` True+ outcome <- OpsRebuild.runCommand (opsEnv False store) emptyCatalogOperations command+ case outcome of+ Failed detail ->+ detail `shouldSatisfy` Text.isPrefixOf "stream-reprojection-group-unregistered:"+ other -> expectationFailure ("expected a typed targeted-repair refusal, got " <> show other)++ it "keeps every typed refusal code distinct and stable" $ \_ -> do+ let group = opsGroupId+ otherGroup = opsGroupBId+ projection = catalogIdentity Catalog.mkProjectionId "ops-projection"+ revision = catalogIdentity Catalog.mkProjectionRevisionId "ops-revision"+ source = catalogIdentity Catalog.mkSourceId "ops-source"+ target = catalogIdentity Catalog.mkTargetId "ops-target"+ dedup = catalogIdentity Catalog.mkDedupKeyId "ops-dedup"+ stream = StreamName "ops-1"+ version = StreamVersion 1+ errors =+ [ Rebuild.StreamReprojectionInvalidPageSize 0,+ Rebuild.StreamReprojectionInvalidMaxEvents 0,+ Rebuild.StreamReprojectionEventLimitExceeded stream 2 1,+ Rebuild.StreamReprojectionGroupUnregistered group,+ Rebuild.StreamReprojectionActiveRebuild group (opsRebuildRunId "ops-active"),+ Rebuild.StreamReprojectionGroupUnavailable group "failed" False False,+ Rebuild.StreamReprojectionSliceDrift group "expected" "actual",+ Rebuild.StreamReprojectionServingRevisionUnavailable group revision,+ Rebuild.StreamReprojectionServingBindingInvalid group revision "invalid binding",+ Rebuild.StreamReprojectionUnknownProjection projection,+ Rebuild.StreamReprojectionProjectionGroupMismatch projection group otherGroup,+ Rebuild.StreamReprojectionPolicyUnavailable revision projection,+ Rebuild.StreamReprojectionSourceMismatch source stream,+ Rebuild.StreamReprojectionHistoryUnavailable (StreamHistoryNotFound stream),+ Rebuild.StreamReprojectionSoftDeleted stream,+ Rebuild.StreamReprojectionTruncated stream version,+ Rebuild.StreamReprojectionForeignEvent stream version,+ Rebuild.StreamReprojectionClearFailed "clear failed",+ Rebuild.StreamReprojectionClearEvidenceInvalid [target] [],+ Rebuild.StreamReprojectionDecodeFailed version (Catalog.ReplayDecodeError "decode failed"),+ Rebuild.StreamReprojectionVerificationFailed "verification failed",+ Rebuild.StreamReprojectionDedupIdentityUnavailable dedup,+ Rebuild.StreamReprojectionHistoryIncomplete version (StreamVersion 0)+ ]+ map OpsRebuild.streamReprojectionErrorCode errors+ `shouldBe` [ "stream-reprojection-invalid-page-size",+ "stream-reprojection-invalid-max-events",+ "stream-reprojection-event-limit-exceeded",+ "stream-reprojection-group-unregistered",+ "stream-reprojection-active-rebuild",+ "stream-reprojection-group-unavailable",+ "stream-reprojection-slice-drift",+ "stream-reprojection-serving-revision-unavailable",+ "stream-reprojection-serving-binding-invalid",+ "stream-reprojection-unknown-projection",+ "stream-reprojection-projection-group-mismatch",+ "stream-reprojection-policy-unavailable",+ "stream-reprojection-source-mismatch",+ "stream-reprojection-history-unavailable",+ "stream-reprojection-soft-deleted",+ "stream-reprojection-truncated",+ "stream-reprojection-foreign-event",+ "stream-reprojection-clear-failed",+ "stream-reprojection-clear-evidence-invalid",+ "stream-reprojection-decode-failed",+ "stream-reprojection-verification-failed",+ "stream-reprojection-dedup-identity-unavailable",+ "stream-reprojection-history-incomplete"+ ]++ describe "catalog rebuild adoption" $ around (withFreshStore fixture) do+ it "previews exact slice changes and adopts them only with force" $ \store -> do+ expectStore store $ runTransaction $ Tx.sql (ByteString.pack "CREATE SCHEMA app; CREATE TABLE app.ops_catalog (id bigint PRIMARY KEY)")+ current <- expectValidatedCatalog (opsCatalog "ops-codec-v1")+ changed <- expectValidatedCatalog (opsCatalog "ops-codec-v2")+ registered <- expectStore store (Rebuild.registerProjectionCatalog current)+ registered `shouldSatisfy` isRight+ let operations = CatalogOperations.projectionCatalogOperations changed+ command = OpsRebuild.Adopt (OpsRebuild.AdoptOptions (NonEmpty.singleton opsGroupId))+ previewEnv =+ OpsEnv+ { store,+ outputMode = HumanTable,+ force = False,+ schemaDrift = [],+ allowSchemaDrift = False+ }++ preview <- OpsRebuild.runCommand previewEnv operations command+ case preview of+ PreviewRequired result invocation -> do+ result.headers `shouldBe` ["name", "kind", "state", "scope", "stored", "current"]+ case result.rows of+ [group, "group", state, scope, stored, currentSlice] : _ -> do+ group `shouldBe` "ops-group"+ state `shouldBe` "slice-changed"+ scope `shouldBe` "adopt"+ stored `shouldSatisfy` Text.isPrefixOf "slice-v6:"+ currentSlice `shouldSatisfy` Text.isPrefixOf "slice-v6:"+ stored `shouldNotBe` currentSlice+ otherRows -> expectationFailure ("unexpected adoption preview rows: " <> show otherRows)+ renderHuman result+ `shouldSatisfy` Text.isInfixOf "adoption changes only keiro-owned registration metadata"+ invocation+ `shouldBe` "'keiro-ops' 'rebuild' 'adopt' 'ops-group' '--force'"+ other -> expectationFailure ("expected adoption preview, got " <> show other)++ applied <- OpsRebuild.runCommand (opsEnv True store) operations command+ case applied of+ Succeeded result -> do+ result.headers `shouldBe` ["name", "kind", "outcome", "detail"]+ result.rows `shouldSatisfy` any (\row -> take 3 row == ["ops-group", "group", "live"])+ case result.jsonValue of+ Aeson.Object fields ->+ KeyMap.lookup "schema" fields+ `shouldBe` Just (Aeson.String "keiro/catalog-adoption-outcome/v2")+ value -> expectationFailure ("expected adoption outcome JSON object, got " <> show value)+ other -> expectationFailure ("expected adoption outcome, got " <> show other)+ registeredChanged <- expectStore store (Rebuild.registerProjectionCatalog changed)+ registeredChanged `shouldSatisfy` isRight+ begun <-+ expectStore+ store+ ( Rebuild.beginGroupRebuild+ changed+ opsGroupId+ Rebuild.RebuildRequest+ { rebuildRunId = opsRunId,+ requestedBy = "keiro-ops-test",+ requestReason = "prove adopted slice can rebuild",+ replayFrom = GlobalPosition 0+ }+ )+ begun `shouldSatisfy` isRight++ it "annotates preview scope and warns about out-of-scope drift" $ \store -> do+ expectStore store $ runTransaction $ Tx.sql (ByteString.pack "CREATE SCHEMA app; CREATE TABLE app.ops_catalog (id bigint PRIMARY KEY); CREATE TABLE app.ops_catalog_b (id bigint PRIMARY KEY)")+ current <- expectValidatedCatalog (opsCatalogPair "ops-codec-a-v1" "ops-codec-b-v1")+ changed <- expectValidatedCatalog (opsCatalogPair "ops-codec-a-v2" "ops-codec-b-v2")+ registered <- expectStore store (Rebuild.registerProjectionCatalog current)+ registered `shouldSatisfy` isRight+ let operations = CatalogOperations.projectionCatalogOperations changed+ command = OpsRebuild.Adopt (OpsRebuild.AdoptOptions (NonEmpty.singleton opsGroupId))+ preview <- OpsRebuild.runCommand (opsEnv False store) operations command+ case preview of+ PreviewRequired result _ -> do+ result.headers `shouldBe` ["name", "kind", "state", "scope", "stored", "current"]+ result.rows `shouldSatisfy` any (\row -> take 4 row == ["ops-group", "group", "slice-changed", "adopt"])+ result.rows `shouldSatisfy` any (\row -> take 4 row == ["ops-group-b", "group", "slice-changed", "skip"])+ case result.jsonValue of+ Aeson.Object fields -> do+ KeyMap.lookup "schema" fields+ `shouldBe` Just (Aeson.String "keiro/catalog-adoption-preview/v2")+ KeyMap.lookup "outOfScopeChangedGroups" fields+ `shouldBe` Just (Aeson.toJSON (["ops-group-b"] :: [Text]))+ value -> expectationFailure ("expected adoption preview JSON object, got " <> show value)+ renderHuman result `shouldSatisfy` Text.isInfixOf "out-of-scope"+ renderHuman result `shouldSatisfy` Text.isInfixOf "ops-group-b"+ other -> expectationFailure ("expected scoped adoption preview, got " <> show other)++ it "refuses an adoption preview for a group absent from the catalog" $ \store -> do+ changed <- expectValidatedCatalog (opsCatalog "ops-codec-v2")+ let missingGroup = catalogIdentity Catalog.mkRebuildGroupId "ops-group-missing"+ operations = CatalogOperations.projectionCatalogOperations changed+ command = OpsRebuild.Adopt (OpsRebuild.AdoptOptions (NonEmpty.singleton missingGroup))+ preview <- OpsRebuild.runCommand (opsEnv False store) operations command+ case preview of+ Failed message -> message `shouldSatisfy` Text.isInfixOf "AdoptGroupNotInCatalog"+ other -> expectationFailure ("expected adoption preview refusal, got " <> show other)++ it "recovers a pre-canonical stranded run through preview and force" $ \store -> do+ expectStore store $ runTransaction $ Tx.sql (ByteString.pack "CREATE SCHEMA app; CREATE TABLE app.ops_catalog (id bigint PRIMARY KEY)")+ let strandedRun = opsRebuildRunId "ops-stranded-run"+ freshRun = opsRebuildRunId "ops-recovery-fresh"+ passingCatalog = opsCatalog "ops-codec-v1"+ failingHook =+ Catalog.RebuildVerification+ { verificationId = "ops-pre-canonical-verification",+ verificationVersion = "v1",+ verifyRebuild = pure (Left "fault injected by keiro-ops recovery spec")+ }+ healthy <- expectValidatedCatalog passingCatalog+ faulted <- expectValidatedCatalog (opsCatalogWithVerifications [failingHook] passingCatalog)+ _ <- expectStore store (Rebuild.registerProjectionCatalog faulted)+ initial <-+ expectStore+ store+ ( Rebuild.startCatalogRebuild+ faulted+ opsGroupId+ ( Rebuild.defaultRebuildOptions+ Rebuild.RebuildRequest+ { rebuildRunId = strandedRun,+ requestedBy = "keiro-ops-test",+ requestReason = "strand a pre-canonical run",+ replayFrom = GlobalPosition 0+ }+ )+ )+ initial `shouldSatisfy` \case+ Left Rebuild.CatalogRebuildVerificationFailed {} -> True+ _ -> False+ expectStore store $ runTransaction $ do+ Tx.sql+ "UPDATE keiro.keiro_projection_rebuild_runs SET group_slice_fingerprint = '$pre-canonical', contract_fingerprint = 'contract-v2:' || repeat('c', 64), runner_format = 'keiro/projection-replay/v2' WHERE run_id = 'ops-stranded-run'"+ Tx.sql+ "UPDATE keiro.keiro_projection_rebuild_groups SET slice_fingerprint = repeat('a', 64) WHERE group_id = 'ops-group'"++ let operations = CatalogOperations.projectionCatalogOperations healthy+ previewEnv =+ OpsEnv+ { store,+ outputMode = HumanTable,+ force = False,+ schemaDrift = [],+ allowSchemaDrift = False+ }+ forceEnv =+ OpsEnv+ { store,+ outputMode = HumanTable,+ force = True,+ schemaDrift = [],+ allowSchemaDrift = False+ }+ abandon =+ OpsRebuild.Abandon+ OpsRebuild.AbandonOptions+ { runId = strandedRun,+ failureCode = "operator.pre-canonical",+ failureDetail = "discard run stranded by migration 0024"+ }++ status <- OpsRebuild.runCommand previewEnv operations (OpsRebuild.Status strandedRun)+ case status of+ Succeeded result -> do+ result.headers+ `shouldBe` ["run", "group", "status", "group_slice", "captured_head", "sources", "adapters", "verifications"]+ result.rows `shouldSatisfy` \case+ [row] -> row !! 3 == "$pre-canonical"+ _ -> False+ other -> expectationFailure ("expected sentinel status, got " <> show other)++ renamedStatus <- OpsRebuild.runCommand previewEnv emptyCatalogOperations (OpsRebuild.Status strandedRun)+ renamedStatus `shouldSatisfy` isSucceeded++ abandonPreview <- OpsRebuild.runCommand previewEnv operations abandon+ case abandonPreview of+ PreviewRequired result invocation -> do+ result.rows `shouldSatisfy` \case+ [row] -> row !! 3 == "$pre-canonical"+ _ -> False+ invocation `shouldSatisfy` Text.isSuffixOf "'--force'"+ other -> expectationFailure ("expected abandon preview, got " <> show other)++ abandoned <- OpsRebuild.runCommand forceEnv operations abandon+ case abandoned of+ Succeeded result ->+ result.rows `shouldSatisfy` \case+ [row] -> row !! 2 == "RebuildRunFailed"+ _ -> False+ other -> expectationFailure ("expected forced abandon, got " <> show other)++ let adopt = OpsRebuild.Adopt (OpsRebuild.AdoptOptions (NonEmpty.singleton opsGroupId))+ adoptionPreview <- OpsRebuild.runCommand previewEnv operations adopt+ case adoptionPreview of+ PreviewRequired result _ ->+ result.rows `shouldSatisfy` \case+ row : _ -> row !! 1 == "group" && row !! 2 == "stale-format"+ _ -> False+ other -> expectationFailure ("expected adoption preview, got " <> show other)+ adopted <- OpsRebuild.runCommand forceEnv operations adopt+ case adopted of+ Succeeded result ->+ result.rows `shouldSatisfy` \case+ row : _ -> row !! 1 == "group" && row !! 2 == "failed" && Text.isPrefixOf "slice-v6:" (row !! 3)+ _ -> False+ other -> expectationFailure ("expected forced adoption, got " <> show other)++ promoted <-+ OpsRebuild.runCommand+ forceEnv+ operations+ ( OpsRebuild.Start+ OpsRebuild.StartOptions+ { groupId = opsGroupId,+ runId = freshRun,+ requestedBy = "keiro-ops-test",+ reason = "fresh canonical recovery run",+ replayFrom = GlobalPosition 0,+ pageSize = 100+ }+ )+ case promoted of+ Succeeded result ->+ result.rows `shouldSatisfy` \case+ [row] -> row !! 2 == "RebuildRunPromoted"+ _ -> False+ other -> expectationFailure ("expected promoted fresh run, got " <> show other)++ describe "durable checkpoint inventory" do+ it "mounts both read-only commands in the standalone tree without a lag alias" do+ isParseSuccess (parseOps Ops.emptyAppHooks ["stream", "subscriptions"]) `shouldBe` True+ isParseSuccess (parseOps Ops.emptyAppHooks ["projection", "position", "--subscription", "orders"]) `shouldBe` True+ isParseFailure (parseOps Ops.emptyAppHooks ["projection", "lag", "--subscription", "orders"]) `shouldBe` True+ OpsStream.isMutation OpsStream.Subscriptions `shouldBe` False+ OpsProjection.isMutation (OpsProjection.Position "orders") `shouldBe` False++ around (withFreshStore fixture) do+ it "returns the captured store position with empty durable rows and null summaries" $ \store -> do+ streamOutcome <- OpsStream.runCommand (opsEnv False store) OpsStream.Subscriptions+ Succeeded streamResult <- pure streamOutcome+ streamResult.rows `shouldBe` []+ streamResult.jsonValue+ `shouldBe` object+ [ "store_position" .= (0 :: Int),+ "visible_store_head" .= (0 :: Int),+ "checkpoints" .= ([] :: [Aeson.Value])+ ]++ projectionOutcome <- OpsProjection.runCommand (opsEnv False store) (OpsProjection.Position "missing")+ Succeeded projectionResult <- pure projectionOutcome+ projectionResult.rows+ `shouldBe` [["missing", "", "", "", "0", "0", "", "", ""]]+ projectionResult.jsonValue+ `shouldBe` object+ [ "subscription" .= ("missing" :: Text),+ "store_position" .= (0 :: Int),+ "visible_store_head" .= (0 :: Int),+ "members" .= ([] :: [Aeson.Value]),+ "minimum_checkpoint_position" .= (Nothing :: Maybe Int64),+ "maximum_global_position_distance" .= (Nothing :: Maybe Int64)+ ]++ it "lists stopped-worker rows in name/member order and derives the member-aware floor" $ \store -> do+ seedCheckpointInventory store++ streamOutcome <- OpsStream.runCommand (opsEnv False store) OpsStream.Subscriptions+ Succeeded streamResult <- pure streamOutcome+ streamResult.rows+ `shouldBe` [ ["billing", "0", "4", "2026-08-09T14:02:00Z", "5", "5", "1"],+ ["orders", "0", "2", "2026-08-09T14:00:00Z", "5", "5", "3"],+ ["orders", "1", "3", "2026-08-09T14:01:00Z", "5", "5", "2"]+ ]+ streamResult.jsonValue+ `shouldBe` object+ [ "store_position" .= (5 :: Int),+ "visible_store_head" .= (5 :: Int),+ "checkpoints"+ .= [ checkpointJsonFixture "billing" 0 4 "2026-08-09T14:02:00Z" 1,+ checkpointJsonFixture "orders" 0 2 "2026-08-09T14:00:00Z" 3,+ checkpointJsonFixture "orders" 1 3 "2026-08-09T14:01:00Z" 2+ ]+ ]++ projectionOutcome <- OpsProjection.runCommand (opsEnv False store) (OpsProjection.Position "orders")+ Succeeded projectionResult <- pure projectionOutcome+ projectionResult.rows+ `shouldBe` [ ["orders", "0", "2", "2026-08-09T14:00:00Z", "5", "5", "3", "2", "3"],+ ["orders", "1", "3", "2026-08-09T14:01:00Z", "5", "5", "2", "2", "3"]+ ]+ projectionResult.jsonValue+ `shouldBe` object+ [ "subscription" .= ("orders" :: Text),+ "store_position" .= (5 :: Int),+ "visible_store_head" .= (5 :: Int),+ "members"+ .= [ checkpointJsonFixture "orders" 0 2 "2026-08-09T14:00:00Z" 3,+ checkpointJsonFixture "orders" 1 3 "2026-08-09T14:01:00Z" 2+ ],+ "minimum_checkpoint_position" .= (2 :: Int),+ "maximum_global_position_distance" .= (3 :: Int)+ ]++ it "diverges store_position from visible_store_head after a hard delete" $ \store -> do+ seedCheckpointInventory store+ Just _ <- expectStore store (hardDeleteStream (StreamName "checkpoint-inventory-5"))++ streamOutcome <- OpsStream.runCommand (opsEnv False store) OpsStream.Subscriptions+ Succeeded streamResult <- pure streamOutcome+ streamResult.rows+ `shouldBe` [ ["billing", "0", "4", "2026-08-09T14:02:00Z", "5", "4", "0"],+ ["orders", "0", "2", "2026-08-09T14:00:00Z", "5", "4", "2"],+ ["orders", "1", "3", "2026-08-09T14:01:00Z", "5", "4", "1"]+ ]+ streamResult.jsonValue+ `shouldBe` object+ [ "store_position" .= (5 :: Int),+ "visible_store_head" .= (4 :: Int),+ "checkpoints"+ .= [ checkpointJsonFixture "billing" 0 4 "2026-08-09T14:02:00Z" 0,+ checkpointJsonFixture "orders" 0 2 "2026-08-09T14:00:00Z" 2,+ checkpointJsonFixture "orders" 1 3 "2026-08-09T14:01:00Z" 1+ ]+ ]++ projectionOutcome <- OpsProjection.runCommand (opsEnv False store) (OpsProjection.Position "orders")+ Succeeded projectionResult <- pure projectionOutcome+ projectionResult.rows+ `shouldBe` [ ["orders", "0", "2", "2026-08-09T14:00:00Z", "5", "4", "2", "2", "2"],+ ["orders", "1", "3", "2026-08-09T14:01:00Z", "5", "4", "1", "2", "2"]+ ]+ projectionResult.jsonValue+ `shouldBe` object+ [ "subscription" .= ("orders" :: Text),+ "store_position" .= (5 :: Int),+ "visible_store_head" .= (4 :: Int),+ "members"+ .= [ checkpointJsonFixture "orders" 0 2 "2026-08-09T14:00:00Z" 2,+ checkpointJsonFixture "orders" 1 3 "2026-08-09T14:01:00Z" 1+ ],+ "minimum_checkpoint_position" .= (2 :: Int),+ "maximum_global_position_distance" .= (2 :: Int)+ ]++ describe "selectConnectionString" do+ it "prefers the explicit option, then the Keiro variable, then DATABASE_URL" do+ selectConnectionString (Just "explicit") (Just "keiro") (Just "database")+ `shouldBe` "explicit"+ selectConnectionString Nothing (Just "keiro") (Just "database")+ `shouldBe` "keiro"+ selectConnectionString Nothing Nothing (Just "database")+ `shouldBe` "database"++ it "uses an empty libpq string for standard PG environment fallbacks" do+ selectConnectionString Nothing Nothing Nothing `shouldBe` ""++ describe "parseDuration" do+ it "rejects every non-finite spelling Read Double accepts" do+ parseDuration "NaN"+ `shouldBe` Left "invalid duration \"NaN\": expected a finite, non-negative number of seconds, optionally with an s, m, h, or d suffix"+ parseDuration "-NaN" `shouldSatisfy` isLeft+ parseDuration "Infinity" `shouldSatisfy` isLeft+ parseDuration "-Infinity" `shouldSatisfy` isLeft+ parseDuration "NaNs" `shouldSatisfy` isLeft+ parseDuration "NaNm" `shouldSatisfy` isLeft+ parseDuration "NaNh" `shouldSatisfy` isLeft+ parseDuration "NaNd" `shouldSatisfy` isLeft+ parseDuration "Infinityd" `shouldSatisfy` isLeft++ it "rejects finite durations the timestamptz wire encoding cannot represent" do+ parseDuration "1e13"+ `shouldBe` Left "invalid duration \"1e13\": exceeds the maximum supported duration of 9.0e12 seconds (about 285000 years)"+ parseDuration "1e308" `shouldSatisfy` isLeft+ parseDuration "1e308d" `shouldSatisfy` isLeft+ parseDuration "115740741000000d" `shouldSatisfy` isLeft++ it "still rejects lowercase non-finite spellings, negatives, and junk" do+ parseDuration "nan" `shouldSatisfy` isLeft+ parseDuration "infinity" `shouldSatisfy` isLeft+ parseDuration "-1" `shouldSatisfy` isLeft+ parseDuration "-1s" `shouldSatisfy` isLeft+ parseDuration "soon" `shouldSatisfy` isLeft+ parseDuration "" `shouldSatisfy` isLeft++ it "accepts integers, decimals, scientific notation, and suffixes unchanged" do+ parseDuration "0" `shouldBe` Right 0+ parseDuration "1.5" `shouldBe` Right 1.5+ parseDuration "2592000" `shouldBe` Right 2592000+ parseDuration "1e6" `shouldBe` Right 1000000+ parseDuration "2m" `shouldBe` Right 120+ parseDuration "3h" `shouldBe` Right 10800+ parseDuration "30d" `shouldBe` Right 2592000+ parseDuration "9.0e12" `shouldBe` Right 9000000000000++ describe "renderHuman" do+ it "aligns columns without changing the structured JSON value" do+ let result =+ OpsResult+ { headers = ["name", "status"],+ rows = [["short", "running"], ["longer", "failed"]],+ jsonValue = object ["items" .= (["unchanged"] :: [String])]+ }+ renderHuman result+ `shouldBe` "name status \n------ -------\nshort running\nlonger failed \n"++ describe "keiro-ops numeric argument rejection" do+ it "refuses a NaN duration before any preview or database contact" do+ executable <- keiroOpsExecutable+ (exit, _, errText) <-+ readProcessWithExitCode+ executable+ [ "--database-url",+ "postgresql://nobody@127.0.0.1:1/unreachable",+ "outbox",+ "gc-sent",+ "--older-than",+ "NaN"+ ]+ ""+ exit `shouldBe` ExitFailure 2+ errText `shouldSatisfy` Text.isInfixOf "invalid duration \"NaN\"" . Text.pack+ errText `shouldSatisfy` not . Text.isInfixOf "preview only" . Text.pack+ errText `shouldSatisfy` not . Text.isInfixOf "schema verification" . Text.pack++ describe "keiro-ops executable" $ around (withFreshDatabase fixture) do+ it "emits parseable JSON and refuses a mutation after schema drift" $ \connectionString -> do+ executable <- keiroOpsExecutable+ (listExit, listOutput, listError) <-+ readProcessWithExitCode+ executable+ ["--database-url", Text.unpack connectionString, "wf", "list", "--json"]+ ""+ listExit `shouldBe` ExitSuccess+ listError `shouldBe` ""+ Aeson.eitherDecodeStrict' (Text.Encoding.encodeUtf8 (Text.pack listOutput))+ `shouldBe` Right (Aeson.Array mempty)++ (previewExit, _, previewError) <-+ readProcessWithExitCode+ executable+ [ "--database-url",+ Text.unpack connectionString,+ "wf",+ "gc",+ "run-once",+ "--retention",+ "0s",+ "--json"+ ]+ ""+ previewExit `shouldBe` ExitFailure 1+ previewError `shouldSatisfy` Text.isInfixOf "preview only" . Text.pack+ previewError `shouldSatisfy` not . Text.isInfixOf "keiro-ops: ExitFailure" . Text.pack++ executeSql connectionString "ALTER TABLE keiro.keiro_timers ADD COLUMN ops_test_drift text"+ (mutationExit, _, mutationError) <-+ readProcessWithExitCode+ executable+ [ "--database-url",+ Text.unpack connectionString,+ "wf",+ "gc",+ "run-once",+ "--retention",+ "0s",+ "--force",+ "--json"+ ]+ ""+ mutationExit `shouldBe` ExitFailure 1+ mutationError `shouldSatisfy` Text.isInfixOf "refusing mutation" . Text.pack++ describe "workflow handlers" $ around (withFreshStore fixture) do+ it "previews and runs one bounded application-registry resume pass" $ \store -> do+ let ref = OpsWorkflow.WorkflowRef "approval" "wf-resume"+ registry = Map.singleton (WorkflowName "approval") (WorkflowDef (\_ -> pure ("done" :: Text)))+ hook = Just (registry, defaultWorkflowResumeOptions)+ command = OpsWorkflow.ResumeOnce (OpsWorkflow.ResumeOptions 1)+ seedStep store ref "received" Aeson.Null++ preview <- OpsWorkflow.runCommandWithResume hook (opsEnv False store) command+ preview `shouldSatisfy` isPreview+ workflowStatus store ref `shouldReturn` Just Instance.WfRunning++ applied <- OpsWorkflow.runCommandWithResume hook (opsEnv True store) command+ applied `shouldSatisfy` isSucceeded+ jsonInteger "completed" applied `shouldBe` Just 1+ jsonInteger "advanced" applied `shouldBe` Just 1+ jsonStringArray "unregistered_names" applied `shouldBe` Just []+ workflowStatus store ref `shouldReturn` Just Instance.WfCompleted++ it "reports advanced work and the exact unregistered workflow names" $ \store -> do+ let registered = OpsWorkflow.WorkflowRef "approval" "wf-resume-registered"+ unregistered = OpsWorkflow.WorkflowRef "retired-approval" "wf-resume-unregistered"+ registry = Map.singleton (WorkflowName "approval") (WorkflowDef (\_ -> pure ("done" :: Text)))+ hook = Just (registry, defaultWorkflowResumeOptions)+ command = OpsWorkflow.ResumeOnce (OpsWorkflow.ResumeOptions 2)+ seedStep store registered "received" Aeson.Null+ seedStep store unregistered "received" Aeson.Null++ applied <- OpsWorkflow.runCommandWithResume hook (opsEnv True store) command+ applied `shouldSatisfy` isSucceeded+ jsonInteger "discovered" applied `shouldBe` Just 2+ jsonInteger "advanced" applied `shouldBe` Just 1+ jsonInteger "unknown_name" applied `shouldBe` Just 1+ jsonStringArray "unregistered_names" applied `shouldBe` Just ["retired-approval"]+ workflowStatus store registered `shouldReturn` Just Instance.WfCompleted+ workflowStatus store unregistered `shouldReturn` Just Instance.WfRunning++ it "classifies a due sleep with no timer worker as blocked, not advanced" $ \store -> do+ let ref = OpsWorkflow.WorkflowRef "approval" "wf-resume-due-sleep"+ registry =+ Map.singleton+ (WorkflowName "approval")+ (WorkflowDef (\_ -> sleepNamed (StepName "wait") (-1) *> pure ("done" :: Text)))+ hook = Just (registry, defaultWorkflowResumeOptions)+ command = OpsWorkflow.ResumeOnce (OpsWorkflow.ResumeOptions 1)+ seedStep store ref "received" Aeson.Null++ first <- OpsWorkflow.runCommandWithResume hook (opsEnv True store) command+ first `shouldSatisfy` isSucceeded+ jsonInteger "discovered" first `shouldBe` Just 1+ jsonInteger "advanced" first `shouldBe` Just 0+ jsonInteger "still_suspended" first `shouldBe` Just 1+ jsonInteger "sleep_due" first `shouldBe` Just 1+ humanField "sleep_due" first `shouldBe` Just "1"++ second <- OpsWorkflow.runCommandWithResume hook (opsEnv True store) command+ second `shouldSatisfy` isSucceeded+ jsonInteger "discovered" second `shouldBe` Just 1+ jsonInteger "advanced" second `shouldBe` Just 0+ jsonInteger "still_suspended" second `shouldBe` Just 1+ jsonInteger "sleep_due" second `shouldBe` Just 1+ humanField "sleep_due" second `shouldBe` Just "1"++ it "lists and decodes a real journal without mutating it" $ \store -> do+ let ref = OpsWorkflow.WorkflowRef "approval" "wf-1"+ seedStep store ref "received" (object ["amount" .= (42 :: Int)])++ listed <-+ OpsWorkflow.runCommand+ (opsEnv False store)+ (OpsWorkflow.List (OpsWorkflow.ListOptions [] Nothing Nothing 100))+ resultArrayLength listed `shouldBe` Just 1++ journal <-+ OpsWorkflow.runCommand+ (opsEnv False store)+ (OpsWorkflow.Journal (OpsWorkflow.InspectOptions ref Nothing))+ journalEventCount journal `shouldBe` Just 1++ row <- runStoreIO store (Instance.lookupInstance (WorkflowName "approval") (WorkflowId "wf-1"))+ fmap (fmap (.wakeAfter)) row `shouldBe` Right (Just Nothing)++ it "applies exact name/status filters and keyset cursors" $ \store -> do+ let first = OpsWorkflow.WorkflowRef "approval" "wf-a"+ second = OpsWorkflow.WorkflowRef "approval" "wf-b"+ other = OpsWorkflow.WorkflowRef "billing" "wf-c"+ seedStep store first "received" Aeson.Null+ seedStep store second "received" Aeson.Null+ seedStep store other "received" Aeson.Null++ page <-+ OpsWorkflow.runCommand+ (opsEnv False store)+ ( OpsWorkflow.List+ ( OpsWorkflow.ListOptions+ [Instance.WfRunning]+ (Just "approval")+ (Just ("approval", "wf-a"))+ 1+ )+ )+ resultArrayLength page `shouldBe` Just 1+ firstWorkflowId page `shouldBe` Just "wf-b"++ it "previews cancellation without mutation, then records it with force" $ \store -> do+ let ref = OpsWorkflow.WorkflowRef "approval" "wf-2"+ seedStep store ref "received" Aeson.Null++ preview <- OpsWorkflow.runCommand (opsEnv False store) (OpsWorkflow.Cancel ref)+ preview `shouldSatisfy` isPreview+ workflowStatus store ref `shouldReturn` Just Instance.WfRunning++ applied <- OpsWorkflow.runCommand (opsEnv True store) (OpsWorkflow.Cancel ref)+ applied `shouldSatisfy` isSucceeded+ workflowStatus store ref `shouldReturn` Just Instance.WfCancelled++ it "previews and applies failed-workflow resurrection and lease release" $ \store -> do+ let ref = OpsWorkflow.WorkflowRef "approval" "wf-recover"+ now <- getCurrentTime+ expectStore store $+ appendJournalEntry+ (WorkflowName ref.workflowName)+ (WorkflowId ref.workflowId)+ WorkflowFailed {reason = "exhausted", recordedAt = now}++ resurrectPreview <- OpsWorkflow.runCommand (opsEnv False store) (OpsWorkflow.Resurrect ref)+ resurrectPreview `shouldSatisfy` isPreview+ workflowStatus store ref `shouldReturn` Just Instance.WfFailed++ resurrected <- OpsWorkflow.runCommand (opsEnv True store) (OpsWorkflow.Resurrect ref)+ resurrected `shouldSatisfy` isSucceeded+ workflowStatus store ref `shouldReturn` Just Instance.WfRunning++ claimed <-+ expectStore store $+ Instance.claimInstance+ "wedged-worker"+ 300+ (WorkflowName ref.workflowName)+ (WorkflowId ref.workflowId)+ claimed `shouldBe` Instance.ClaimAcquired++ releasePreview <- OpsWorkflow.runCommand (opsEnv False store) (OpsWorkflow.ReleaseLease ref)+ releasePreview `shouldSatisfy` isPreview+ workflowLeaseOwner store ref `shouldReturn` Just "wedged-worker"++ released <- OpsWorkflow.runCommand (opsEnv True store) (OpsWorkflow.ReleaseLease ref)+ released `shouldSatisfy` isSucceeded+ workflowLeaseOwner store ref `shouldReturn` Nothing++ it "previews and signals an awakeable through the supported library path" $ \store -> do+ let ref = OpsWorkflow.WorkflowRef "approval" "wf-3"+ awakeableId = maybe (error "test UUID") Function.id (UUID.fromString "018f5f43-8a70-7b9a-9a9b-59d391a76710")+ seedStep store ref "awkid:approval" (Aeson.toJSON (AwakeableId awakeableId))+ expectStore store $ runTransaction (Awakeable.registerAwakeableTx awakeableId "approval" "wf-3")++ preview <-+ OpsWorkflow.runCommand+ (opsEnv False store)+ (OpsWorkflow.Awakeable (OpsWorkflow.AwakeableSignal awakeableId (OpsWorkflow.PayloadArg "{\"approved\":true}" (object ["approved" .= True]))))+ preview `shouldSatisfy` isPreview+ awakeableStatus store awakeableId `shouldReturn` Just Awakeable.Pending++ applied <-+ OpsWorkflow.runCommand+ (opsEnv True store)+ (OpsWorkflow.Awakeable (OpsWorkflow.AwakeableSignal awakeableId (OpsWorkflow.PayloadArg "{\"approved\":true}" (object ["approved" .= True]))))+ applied `shouldSatisfy` isSucceeded+ awakeableStatus store awakeableId `shouldReturn` Just Awakeable.Completed++ it "previews the exact GC candidates before deleting them" $ \store -> do+ let ref = OpsWorkflow.WorkflowRef "approval" "wf-4"+ gcOptions = OpsWorkflow.GcOptions 0 10+ seedStep store ref "received" Aeson.Null+ _ <- OpsWorkflow.runCommand (opsEnv True store) (OpsWorkflow.Cancel ref)++ preview <- OpsWorkflow.runCommand (opsEnv False store) (OpsWorkflow.GcRunOnce gcOptions)+ resultArrayLengthFrom "candidates" preview `shouldBe` Just 1++ applied <- OpsWorkflow.runCommand (opsEnv True store) (OpsWorkflow.GcRunOnce gcOptions)+ applied `shouldSatisfy` isSucceeded+ workflowStatus store ref `shouldReturn` Nothing++ describe "timer handlers" $ around (withFreshStore fixture) do+ it "previews and dispatches one bounded due-timer pass through the mounted hook" $ \store -> do+ now <- getCurrentTime+ let request = timerRequest "018f5f43-8a70-7b9a-9a9b-59d391a76722" (addUTCTime (-60) now)+ fire _ = pure (Just (EventId (testUuid "018f5f43-8a70-7b9a-9a9b-59d391a76723")))+ command = OpsTimer.DrainOnce (OpsTimer.DrainOptions 1)+ expectStore store (runTransaction (Timer.scheduleTimerTx request))++ preview <- OpsTimer.runCommandWithFire (Just fire) (opsEnv False store) command+ preview `shouldSatisfy` isPreview+ timerStatus store request.timerId `shouldReturn` Just Timer.Scheduled++ applied <- OpsTimer.runCommandWithFire (Just fire) (opsEnv True store) command+ applied `shouldSatisfy` isSucceeded+ jsonInteger "processed" applied `shouldBe` Just 1+ timerStatus store request.timerId `shouldReturn` Just Timer.Fired++ it "lists, previews, requeues, and dead-letters a stuck timer" $ \store -> do+ now <- getCurrentTime+ let request = timerRequest "018f5f43-8a70-7b9a-9a9b-59d391a76720" (addUTCTime (-60) now)+ timerId = request.timerId+ expectStore store (runTransaction (Timer.scheduleTimerTx request))+ claimed <- expectStore store (Timer.claimDueTimer now)+ fmap (.status) claimed `shouldBe` Just Timer.Firing++ tooManyAttempts <-+ OpsTimer.runCommand+ (opsEnv False store)+ (OpsTimer.StuckList (OpsTimer.StuckListOptions Nothing (Just 2)))+ resultArrayLength tooManyAttempts `shouldBe` Just 0++ listed <-+ OpsTimer.runCommand+ (opsEnv False store)+ (OpsTimer.StuckList (OpsTimer.StuckListOptions Nothing Nothing))+ resultArrayLength listed `shouldBe` Just 1++ preview <- OpsTimer.runCommand (opsEnv False store) (OpsTimer.Requeue timerId)+ preview `shouldSatisfy` isPreview+ timerStatus store timerId `shouldReturn` Just Timer.Firing++ requeued <- OpsTimer.runCommand (opsEnv True store) (OpsTimer.Requeue timerId)+ requeued `shouldSatisfy` isSucceeded+ timerStatus store timerId `shouldReturn` Just Timer.Scheduled++ retriedClaim <- expectStore store (Timer.claimDueTimer now)+ retriedClaim `shouldSatisfy` isJust+ retried <-+ OpsTimer.runCommand+ (opsEnv False store)+ (OpsTimer.StuckList (OpsTimer.StuckListOptions Nothing (Just 2)))+ resultArrayLength retried `shouldBe` Just 1++ deadPreview <- OpsTimer.runCommand (opsEnv False store) (OpsTimer.DeadLetter timerId "poison payload")+ deadPreview `shouldSatisfy` isPreview+ timerStatus store timerId `shouldReturn` Just Timer.Firing++ dead <- OpsTimer.runCommand (opsEnv True store) (OpsTimer.DeadLetter timerId "poison payload")+ dead `shouldSatisfy` isSucceeded+ timerStatus store timerId `shouldReturn` Just Timer.Dead++ it "previews and cancels a scheduled timer" $ \store -> do+ now <- getCurrentTime+ let request = timerRequest "018f5f43-8a70-7b9a-9a9b-59d391a76721" (addUTCTime 3600 now)+ timerId = request.timerId+ expectStore store (runTransaction (Timer.scheduleTimerTx request))++ preview <- OpsTimer.runCommand (opsEnv False store) (OpsTimer.Cancel timerId)+ preview `shouldSatisfy` isPreview+ timerStatus store timerId `shouldReturn` Just Timer.Scheduled++ cancelled <- OpsTimer.runCommand (opsEnv True store) (OpsTimer.Cancel timerId)+ cancelled `shouldSatisfy` isSucceeded+ timerStatus store timerId `shouldReturn` Just Timer.Cancelled++ describe "outbox handlers" $ around (withFreshStore fixture) do+ it "lists backlog and previews stale recovery without mutation" $ \store -> do+ now <- getCurrentTime+ let outboxId = testOutboxId "018f5f43-8a70-7b9a-9a9b-59d391a76801"+ event = sampleIntegrationEvent now "outbox-message"+ expectStore store (runTransaction (Outbox.enqueueOutboxTx (Outbox.OutboxMessage outboxId event)))++ backlog <- OpsOutbox.runCommand (opsEnv False store) OpsOutbox.Backlog+ resultCount backlog `shouldBe` Just 1++ claimNow <- getCurrentTime+ _ <- expectStore store (Outbox.claimOutboxBatch Outbox.BestEffort 1 claimNow)+ preview <- OpsOutbox.runCommand (opsEnv False store) (OpsOutbox.RequeueStuck 0 10)+ preview `shouldSatisfy` isPreview+ outboxStatus store outboxId `shouldReturn` Just Outbox.OutboxPublishing++ applied <- OpsOutbox.runCommand (opsEnv True store) (OpsOutbox.RequeueStuck 0 10)+ applied `shouldSatisfy` isSucceeded+ outboxStatus store outboxId `shouldReturn` Just Outbox.OutboxFailed++ it "surfaces dispatch dead letters through the supported API" $ \store -> do+ let sourceEvent = EventId (testUuid "018f5f43-8a70-7b9a-9a9b-59d391a76802")+ expectStore store $+ recordDispatchDeadLetter+ DispatchDeadLetter+ { dispatcherKind = DispatcherProcessManager,+ dispatcherName = "ops-pm",+ correlationId = "order-1",+ sourceEventId = sourceEvent,+ sourceGlobalPosition = GlobalPosition 1,+ emitIndex = 0,+ targetStreamName = StreamName "order-1",+ errorClass = "rejected",+ errorDetail = "operator fixture",+ attemptCount = 1+ }+ listed <- OpsOutbox.runCommand (opsEnv False store) (OpsOutbox.DispatchDeadLetters "ops-pm" 10)+ resultArrayLength listed `shouldBe` Just 1++ describe "inbox handlers" $ around (withFreshStore fixture) do+ it "previews poison marking and GC without bypassing inbox APIs" $ \store -> do+ now <- getCurrentTime+ let poison = sampleIntegrationEvent now "poison-message"+ completed = sampleIntegrationEvent now "completed-message"+ seedInbox store poison+ seedInbox store completed++ preview <- OpsInbox.runCommand (opsEnv False store) (OpsInbox.MarkFailed poison.source poison.messageId "poison")+ preview `shouldSatisfy` isPreview+ inboxStatus store poison.source poison.messageId `shouldReturn` Just Inbox.InboxCompleted++ marked <- OpsInbox.runCommand (opsEnv True store) (OpsInbox.MarkFailed poison.source poison.messageId "poison")+ marked `shouldSatisfy` isSucceeded+ inboxStatus store poison.source poison.messageId `shouldReturn` Just Inbox.InboxFailed++ gcPreview <- OpsInbox.runCommand (opsEnv False store) (OpsInbox.Gc 0)+ gcPreview `shouldSatisfy` isPreview+ inboxStatus store completed.source completed.messageId `shouldReturn` Just Inbox.InboxCompleted++ gcApplied <- OpsInbox.runCommand (opsEnv True store) (OpsInbox.Gc 0)+ gcApplied `shouldSatisfy` isSucceeded+ inboxStatus store completed.source completed.messageId `shouldReturn` Nothing++ describe "pgmq handlers" $ around (withFreshStore fixture) do+ it "previews and redrives a DLQ entry, which is then consumable" $ \store -> do+ let queue = "keiro_ops_test.redrive"+ job = rawValueJob queue+ runPgmqUnit action = do+ result <- runJobEff (JobRuntime store.pool Nothing) action+ either (fail . show) pure result+ depths = do+ result <- runJobEff (JobRuntime store.pool Nothing) $ do+ mainMetrics <- jobQueueMetrics job+ dlqMetrics <- jobDlqMetrics job+ pure (mainMetrics.queueLength, dlqMetrics.queueLength)+ either (fail . show) pure result+ runPgmqUnit $ do+ ensureJobQueue job+ _ <- enqueue job (object ["kind" .= ("poison" :: Text)])+ _ <- runJobOnce 1 job (\_ -> pure (Dead "bad"))+ pure ()++ preview <- OpsPgmq.runCommand (opsEnv False store) (OpsPgmq.Dlq (OpsPgmq.Redrive queue 10))+ preview `shouldSatisfy` isPreview+ (mainBefore, dlqBefore) <- depths+ (mainBefore, dlqBefore) `shouldBe` (0, 1)++ applied <- OpsPgmq.runCommand (opsEnv True store) (OpsPgmq.Dlq (OpsPgmq.Redrive queue 10))+ applied `shouldSatisfy` isSucceeded+ (mainAfter, dlqAfter) <- depths+ (mainAfter, dlqAfter) `shouldBe` (1, 0)++ runPgmqUnit (runJobOnce 1 job (\_ -> pure Done))+ (mainFinal, _) <- depths+ mainFinal `shouldBe` 0++ runPgmqUnit $ do+ _ <- enqueue job (object ["kind" .= ("purge-me" :: Text)])+ _ <- runJobOnce 1 job (\_ -> pure (Dead "still bad"))+ pure ()+ purgePreview <- OpsPgmq.runCommand (opsEnv False store) (OpsPgmq.Dlq (OpsPgmq.Purge queue))+ purgePreview `shouldSatisfy` isPreview+ (_, dlqBeforePurge) <- depths+ dlqBeforePurge `shouldBe` 1++ purged <- OpsPgmq.runCommand (opsEnv True store) (OpsPgmq.Dlq (OpsPgmq.Purge queue))+ purged `shouldSatisfy` isSucceeded+ (_, dlqAfterPurge) <- depths+ dlqAfterPurge `shouldBe` 0++ describe "projection handlers" $ around (withFreshStore fixture) do+ it "prunes only the named dedup rows" $ \store -> do+ _ <- seedKirokuEvent store "projection-source" "018f5f43-8a70-7b9a-9a9b-59d391a76810" Nothing+ events <- expectStore store (readStreamForward (StreamName "projection-source") (StreamVersion 0) 1)+ let recorded = Vector.head events+ projection =+ Projection.AsyncProjection+ { name = "ops-dedup",+ readModelName = "ops-read-model",+ subscriptionName = "ops-projection",+ applyRecorded = \_ -> pure (),+ idempotencyKey = (.eventId)+ }+ _ <- expectStore store (runTransaction (Projection.applyAsyncProjectionUnfenced projection recorded))+ future <- addUTCTime 60 <$> getCurrentTime+ prunePreview <- OpsProjection.runCommand (opsEnv False store) (OpsProjection.PruneDedup "ops-dedup" future)+ prunePreview `shouldSatisfy` isPreview+ jsonIntegerFromPreview "affected" prunePreview `shouldBe` Just 1+ pruned <- OpsProjection.runCommand (opsEnv True store) (OpsProjection.PruneDedup "ops-dedup" future)+ jsonInteger "affected" pruned `shouldBe` Just 1++ describe "shard handlers" $ around (withFreshStore fixture) do+ it "previews exact buckets and relinquishes them for another worker" $ \store -> do+ let subscription = SubscriptionName "ops-shards"+ worker = Shard.WorkerId (testUuid "018f5f43-8a70-7b9a-9a9b-59d391a76803")+ lease = Shard.ShardLease subscription worker 2 300+ expectStore store (Shard.ensureShards lease)+ _ <- expectStore store (Shard.acquireOwnedBuckets lease 1)+ _ <- expectStore store (Shard.acquireOwnedBuckets lease 1)++ status <- OpsShard.runCommand (opsEnv False store) (OpsShard.Status "ops-shards")+ resultArrayLengthFromObject "ownership" status `shouldBe` Just 2++ preview <- OpsShard.runCommand (opsEnv False store) (OpsShard.Relinquish "ops-shards" worker)+ preview `shouldSatisfy` isPreview+ ownersBefore <- expectStore store (Shard.ownershipSnapshotFor subscription)+ length [() | (_, Just owner, _) <- ownersBefore, owner == worker] `shouldBe` 2++ released <- OpsShard.runCommand (opsEnv True store) (OpsShard.Relinquish "ops-shards" worker)+ released `shouldSatisfy` isSucceeded+ ownersAfter <- expectStore store (Shard.ownershipSnapshotFor subscription)+ ownersAfter `shouldSatisfy` all (\(_, owner, _) -> owner == Nothing)++ let replacement = Shard.WorkerId (testUuid "018f5f43-8a70-7b9a-9a9b-59d391a76804")+ replacementLease = Shard.ShardLease subscription replacement 2 300+ _ <- expectStore store (Shard.acquireOwnedBuckets replacementLease 1)+ _ <- expectStore store (Shard.acquireOwnedBuckets replacementLease 1)+ replacementOwners <- expectStore store (Shard.ownershipSnapshotFor subscription)+ replacementOwners `shouldSatisfy` all (\(_, owner, _) -> owner == Just replacement)++ describe "snapshot handlers" $ around (withFreshStore fixture) do+ it "refuses uncovered truncation, passes matching coverage, and deletes advisories" $ \store -> do+ appended <- seedKirokuEvent store "snapshot-ops" "018f5f43-8a70-7b9a-9a9b-59d391a76811" Nothing+ let expected = OpsSnapshot.ExpectedDiscriminators 7 "regs-v7" "fold-v7"+ expectStore store $+ writeSnapshotRow+ SnapshotWrite+ { streamId = appended.streamId,+ streamVersion = appended.streamVersion,+ state = object ["count" .= (1 :: Int)],+ stateCodecVersion = expected.stateCodecVersion,+ regfileShapeHash = expected.regfileShapeHash,+ stateShapeHash = expected.stateShapeHash+ }++ missing <- OpsSnapshot.runCommand (opsEnv False store) (OpsSnapshot.TruncationPreflight "no-snapshot" (StreamVersion 2) (Just expected))+ jsonBool "passed" missing `shouldBe` Just False++ covered <- OpsSnapshot.runCommand (opsEnv False store) (OpsSnapshot.TruncationPreflight "snapshot-ops" (StreamVersion 2) (Just expected))+ jsonBool "passed" covered `shouldBe` Just True++ preview <- OpsSnapshot.runCommand (opsEnv False store) (OpsSnapshot.Delete "snapshot-ops")+ preview `shouldSatisfy` isPreview+ beforeDelete <- expectStore store (lookupSnapshotRow appended.streamId)+ beforeDelete `shouldSatisfy` isJust++ deleted <- OpsSnapshot.runCommand (opsEnv True store) (OpsSnapshot.Delete "snapshot-ops")+ deleted `shouldSatisfy` isSucceeded+ expectStore store (lookupSnapshotRow appended.streamId) `shouldReturn` Nothing++ describe "stream handlers" $ around (withFreshStore fixture) do+ it "reads causation and applies reversible lifecycle operations" $ \store -> do+ first <- seedKirokuEvent store "stream-ops" "018f5f43-8a70-7b9a-9a9b-59d391a76812" Nothing+ second <- seedKirokuEvent store "stream-ops" "018f5f43-8a70-7b9a-9a9b-59d391a76813" (Just (eventUuid first))++ shown <- OpsStream.runCommand (opsEnv False store) (OpsStream.Show "stream-ops" (StreamVersion 0) 10)+ resultArrayLengthFromObject "events" shown `shouldBe` Just 2++ causes <- OpsStream.runCommand (opsEnv False store) (OpsStream.Causation (EventId (eventUuid second)))+ resultArrayLength causes `shouldBe` Just 2++ softPreview <- OpsStream.runCommand (opsEnv False store) (OpsStream.SoftDelete "stream-ops")+ softPreview `shouldSatisfy` isPreview+ streamDeleted store "stream-ops" `shouldReturn` Just False++ softDeleted <- OpsStream.runCommand (opsEnv True store) (OpsStream.SoftDelete "stream-ops")+ softDeleted `shouldSatisfy` isSucceeded+ streamDeleted store "stream-ops" `shouldReturn` Just True++ restored <- OpsStream.runCommand (opsEnv True store) (OpsStream.Undelete "stream-ops")+ restored `shouldSatisfy` isSucceeded+ streamDeleted store "stream-ops" `shouldReturn` Just False++ it "previews and applies truncate markers and permanent deletion" $ \store -> do+ _ <- seedKirokuEvent store "stream-destructive" "018f5f43-8a70-7b9a-9a9b-59d391a76814" Nothing+ _ <- seedKirokuEvent store "stream-destructive" "018f5f43-8a70-7b9a-9a9b-59d391a76815" Nothing++ truncatePreview <-+ OpsStream.runCommand+ (opsEnv False store)+ (OpsStream.TruncateBefore (OpsStream.SetTruncateBefore "stream-destructive" (StreamVersion 2) Nothing True))+ truncatePreview `shouldSatisfy` isPreview+ streamTruncateBefore store "stream-destructive" `shouldReturn` Just (StreamVersion 0)++ truncated <-+ OpsStream.runCommand+ (opsEnv True store)+ (OpsStream.TruncateBefore (OpsStream.SetTruncateBefore "stream-destructive" (StreamVersion 2) Nothing True))+ truncated `shouldSatisfy` isSucceeded+ streamTruncateBefore store "stream-destructive" `shouldReturn` Just (StreamVersion 2)++ clearPreview <- OpsStream.runCommand (opsEnv False store) (OpsStream.TruncateBefore (OpsStream.ClearTruncateBefore "stream-destructive"))+ clearPreview `shouldSatisfy` isPreview+ streamTruncateBefore store "stream-destructive" `shouldReturn` Just (StreamVersion 2)++ cleared <- OpsStream.runCommand (opsEnv True store) (OpsStream.TruncateBefore (OpsStream.ClearTruncateBefore "stream-destructive"))+ cleared `shouldSatisfy` isSucceeded+ streamTruncateBefore store "stream-destructive" `shouldReturn` Just (StreamVersion 0)++ deletePreview <- OpsStream.runCommand (opsEnv False store) (OpsStream.HardDelete "stream-destructive")+ deletePreview `shouldSatisfy` isPreview+ beforeDelete <- expectStore store (getStream (StreamName "stream-destructive"))+ beforeDelete `shouldSatisfy` isJust++ deleted <- OpsStream.runCommand (opsEnv True store) (OpsStream.HardDelete "stream-destructive")+ deleted `shouldSatisfy` isSucceeded+ expectStore store (getStream (StreamName "stream-destructive")) `shouldReturn` Nothing++data SeededEvent = SeededEvent+ { streamId :: !StreamId,+ streamVersion :: !StreamVersion,+ globalPosition :: !GlobalPosition,+ eventId :: !EventId+ }++sampleIntegrationEvent :: UTCTime -> Text -> IntegrationEvent+sampleIntegrationEvent now messageId =+ IntegrationEvent+ { messageId,+ source = "ops-source",+ destination = "ops-destination",+ key = Just "entity-1",+ eventType = "ops.event",+ schemaVersion = 1,+ contentType = ApplicationJson,+ schemaReference = Nothing,+ sourceEventId = Nothing,+ sourceGlobalPosition = Nothing,+ payloadBytes = ByteString.pack "{\"ok\":true}",+ occurredAt = now,+ causationId = Nothing,+ correlationId = Nothing,+ traceContext = Nothing,+ attributes = Nothing+ }++seedInbox :: KirokuStore -> IntegrationEvent -> IO ()+seedInbox store event = do+ result <-+ expectStore store $+ Inbox.runInboxTransaction+ Nothing+ Inbox.PreferIntegrationMessageId+ event+ Nothing+ (\_ -> pure ())+ result `shouldBe` Right (Inbox.InboxProcessed ())++outboxStatus :: KirokuStore -> Outbox.OutboxId -> IO (Maybe Outbox.OutboxStatus)+outboxStatus store outboxId =+ expectStore store (Outbox.lookupOutbox outboxId) <&> fmap (.status)++inboxStatus :: KirokuStore -> Text -> Text -> IO (Maybe Inbox.InboxStatus)+inboxStatus store source messageId =+ expectStore store (Inbox.lookupInbox source messageId) <&> fmap (.status)++testOutboxId :: String -> Outbox.OutboxId+testOutboxId = Outbox.OutboxId . testUuid++testUuid :: String -> UUID.UUID+testUuid raw = maybe (error "test UUID") Function.id (UUID.fromString raw)++rawValueJob :: Text -> Job Aeson.Value+rawValueJob name =+ Job+ { jobName = name,+ jobQueue = queueRef name,+ jobCodec = aesonJobCodec,+ jobPolicy = defaultRetryPolicy+ }++seedKirokuEvent :: KirokuStore -> Text -> String -> Maybe UUID.UUID -> IO SeededEvent+seedKirokuEvent store name rawId cause = do+ let eventId = EventId (testUuid rawId)+ appended <-+ expectStore store $+ appendToStream+ (StreamName name)+ AnyVersion+ [ EventData+ { eventId = Just eventId,+ eventType = EventType "ops.event",+ payload = object ["stream" .= name],+ metadata = Nothing,+ causationId = cause,+ correlationId = Nothing+ }+ ]+ pure+ SeededEvent+ { streamId = appended.streamId,+ streamVersion = appended.streamVersion,+ globalPosition = appended.globalPosition,+ eventId+ }++seedCheckpointInventory :: KirokuStore -> IO ()+seedCheckpointInventory store = do+ let seeds =+ [ ("checkpoint-inventory-1", "018f5f43-8a70-7b9a-9a9b-59d391a76821"),+ ("checkpoint-inventory-2", "018f5f43-8a70-7b9a-9a9b-59d391a76822"),+ ("checkpoint-inventory-3", "018f5f43-8a70-7b9a-9a9b-59d391a76823"),+ ("checkpoint-inventory-4", "018f5f43-8a70-7b9a-9a9b-59d391a76824"),+ ("checkpoint-inventory-5", "018f5f43-8a70-7b9a-9a9b-59d391a76825")+ ]+ mapM_ (\(name, eventId) -> seedKirokuEvent store name eventId Nothing) seeds+ expectStore store $+ runTransaction $+ Tx.sql+ "INSERT INTO subscriptions (subscription_name, stream_name, consumer_group_member, consumer_group_size, last_seen, updated_at) VALUES ('orders', '$all', 1, 2, 3, '2026-08-09 14:01:00+00'), ('billing', '$all', 0, 1, 4, '2026-08-09 14:02:00+00'), ('orders', '$all', 0, 2, 2, '2026-08-09 14:00:00+00')"++checkpointJsonFixture :: Text -> Int -> Int -> Text -> Int -> Aeson.Value+checkpointJsonFixture subscription member position updatedAt distance =+ object+ [ "subscription" .= subscription,+ "member" .= member,+ "checkpoint_position" .= position,+ "checkpoint_updated_at" .= updatedAt,+ "global_position_distance" .= distance+ ]++eventUuid :: SeededEvent -> UUID.UUID+eventUuid seeded = case seeded.eventId of EventId value -> value++streamDeleted :: KirokuStore -> Text -> IO (Maybe Bool)+streamDeleted store name =+ expectStore store (getStream (StreamName name)) <&> fmap (isJust . (.deletedAt))++streamTruncateBefore :: KirokuStore -> Text -> IO (Maybe StreamVersion)+streamTruncateBefore store name =+ expectStore store (getStream (StreamName name)) <&> fmap (.truncateBefore)++data OpsCatalogEvent++opsCatalog :: Text -> Catalog.ProjectionCatalog+opsCatalog codecFingerprint =+ Catalog.ProjectionCatalog+ { sources =+ [ Catalog.SourceDeclaration+ { sourceId = opsSourceId,+ sourceScope = Catalog.CategorySource (CategoryName "ops-catalog"),+ codecFingerprint,+ claimSite = catalogIdentity Catalog.mkClaimSite "ops-test:source"+ }+ ],+ targets =+ [ Catalog.TargetDeclaration+ { targetId = opsTargetId,+ qualifiedTable = Catalog.QualifiedTable "app" "ops_catalog",+ resetPolicy = Catalog.ClearBeforeReplay,+ dependsOn = [],+ claimSite = catalogIdentity Catalog.mkClaimSite "ops-test:target"+ }+ ],+ rebuildGroups =+ [ Catalog.RebuildGroupDeclaration+ { rebuildGroupId = opsGroupId,+ orderedTargets = [opsTargetId],+ verificationHooks = [],+ claimSite = catalogIdentity Catalog.mkClaimSite "ops-test:group"+ }+ ],+ projectionRevisions = [],+ externalReadContracts = [],+ subscriptions = [],+ dedupKeys = [],+ queryModels = [],+ projectionSets = [Catalog.SomeProjectionSet opsProjectionSet]+ }++opsCatalogPair :: Text -> Text -> Catalog.ProjectionCatalog+opsCatalogPair firstCodec secondCodec =+ let first = opsCatalog firstCodec+ in first+ { Catalog.sources =+ first.sources+ <> [ Catalog.SourceDeclaration+ { sourceId = opsSourceBId,+ sourceScope = Catalog.CategorySource (CategoryName "ops-catalog-b"),+ codecFingerprint = secondCodec,+ claimSite = catalogIdentity Catalog.mkClaimSite "ops-test:source-b"+ }+ ],+ Catalog.targets =+ first.targets+ <> [ Catalog.TargetDeclaration+ { targetId = opsTargetBId,+ qualifiedTable = Catalog.QualifiedTable "app" "ops_catalog_b",+ resetPolicy = Catalog.ClearBeforeReplay,+ dependsOn = [],+ claimSite = catalogIdentity Catalog.mkClaimSite "ops-test:target-b"+ }+ ],+ Catalog.rebuildGroups =+ first.rebuildGroups+ <> [ Catalog.RebuildGroupDeclaration+ { rebuildGroupId = opsGroupBId,+ orderedTargets = [opsTargetBId],+ verificationHooks = [],+ claimSite = catalogIdentity Catalog.mkClaimSite "ops-test:group-b"+ }+ ],+ Catalog.projectionSets = first.projectionSets <> [Catalog.SomeProjectionSet opsProjectionSetB]+ }++opsCatalogWithVerifications :: [Catalog.RebuildVerification] -> Catalog.ProjectionCatalog -> Catalog.ProjectionCatalog+opsCatalogWithVerifications verifications catalog =+ catalog+ { Catalog.rebuildGroups =+ [ group {Catalog.verificationHooks = verifications}+ | group <- catalog.rebuildGroups+ ]+ }++opsProjectionSet :: Catalog.ProjectionSet OpsCatalogEvent+opsProjectionSet =+ Catalog.ProjectionSet+ { projectionSource = opsSourceId,+ projectionDefinitions = NonEmpty.singleton opsProjectionDefinition,+ claimSite = catalogIdentity Catalog.mkClaimSite "ops-test:set"+ }++opsProjectionDefinition :: Catalog.ProjectionDefinition OpsCatalogEvent+opsProjectionDefinition =+ Catalog.ProjectionDefinition+ { projectionId = catalogIdentity Catalog.mkProjectionId "ops-owner",+ rebuildGroup = opsGroupId,+ ownedTargets = NonEmpty.singleton opsTargetId,+ replayPolicy =+ Catalog.Replayable+ Catalog.ReplayAdapter+ { decodeForReplay = const Catalog.ReplayIrrelevant,+ applyForReplay = \_ _ -> pure ()+ },+ handlers =+ NonEmpty.singleton+ ( Catalog.InlineHandler+ Projection.InlineProjection+ { name = "ops-inline",+ apply = \_ _ -> pure ()+ }+ (catalogIdentity Catalog.mkClaimSite "ops-test:inline-handler")+ ),+ claimSite = catalogIdentity Catalog.mkClaimSite "ops-test:projection"+ }++opsProjectionSetB :: Catalog.ProjectionSet OpsCatalogEvent+opsProjectionSetB =+ Catalog.ProjectionSet+ { projectionSource = opsSourceBId,+ projectionDefinitions = NonEmpty.singleton opsProjectionDefinitionB,+ claimSite = catalogIdentity Catalog.mkClaimSite "ops-test:set-b"+ }++opsProjectionDefinitionB :: Catalog.ProjectionDefinition OpsCatalogEvent+opsProjectionDefinitionB =+ Catalog.ProjectionDefinition+ { projectionId = catalogIdentity Catalog.mkProjectionId "ops-owner-b",+ rebuildGroup = opsGroupBId,+ ownedTargets = NonEmpty.singleton opsTargetBId,+ replayPolicy =+ Catalog.Replayable+ Catalog.ReplayAdapter+ { decodeForReplay = const Catalog.ReplayIrrelevant,+ applyForReplay = \_ _ -> pure ()+ },+ handlers =+ NonEmpty.singleton+ ( Catalog.InlineHandler+ Projection.InlineProjection+ { name = "ops-inline-b",+ apply = \_ _ -> pure ()+ }+ (catalogIdentity Catalog.mkClaimSite "ops-test:inline-handler-b")+ ),+ claimSite = catalogIdentity Catalog.mkClaimSite "ops-test:projection-b"+ }++opsGroupId :: Catalog.RebuildGroupId+opsGroupId = catalogIdentity Catalog.mkRebuildGroupId "ops-group"++opsSourceId :: Catalog.SourceId+opsSourceId = catalogIdentity Catalog.mkSourceId "ops-source"++opsTargetId :: Catalog.TargetId+opsTargetId = catalogIdentity Catalog.mkTargetId "ops-target"++opsGroupBId :: Catalog.RebuildGroupId+opsGroupBId = catalogIdentity Catalog.mkRebuildGroupId "ops-group-b"++opsSourceBId :: Catalog.SourceId+opsSourceBId = catalogIdentity Catalog.mkSourceId "ops-source-b"++opsTargetBId :: Catalog.TargetId+opsTargetBId = catalogIdentity Catalog.mkTargetId "ops-target-b"++opsRunId :: Rebuild.RebuildRunId+opsRunId =+ opsRebuildRunId "ops-adoption-run"++opsRebuildRunId :: Text -> Rebuild.RebuildRunId+opsRebuildRunId identity =+ case Rebuild.mkRebuildRunId identity of+ Left err -> error (Text.unpack err)+ Right value -> value++catalogIdentity :: (Show err) => (Text -> Either err value) -> Text -> value+catalogIdentity constructor value =+ case constructor value of+ Left err -> error (show err)+ Right identity -> identity++expectValidatedCatalog :: Catalog.ProjectionCatalog -> IO Catalog.ValidatedProjectionCatalog+expectValidatedCatalog catalog =+ case Catalog.validateProjectionCatalog catalog of+ Catalog.Failure diagnostics -> expectationFailure (show diagnostics) >> error "unreachable"+ Catalog.Success validated -> pure validated++opsEnv :: Bool -> KirokuStore -> OpsEnv+opsEnv force store =+ OpsEnv+ { store,+ outputMode = Json,+ force,+ schemaDrift = [],+ allowSchemaDrift = False+ }++seedStep :: KirokuStore -> OpsWorkflow.WorkflowRef -> Text -> Aeson.Value -> IO ()+seedStep store ref stepName payload = do+ now <- getCurrentTime+ expectStore store $+ appendJournalEntry+ (WorkflowName ref.workflowName)+ (WorkflowId ref.workflowId)+ StepRecorded {stepName, result = payload, recordedAt = now}++expectStore :: KirokuStore -> Eff '[Store, Error StoreError, IOE] a -> IO a+expectStore store action = runStoreIO store action >>= either (fail . show) pure++workflowStatus :: KirokuStore -> OpsWorkflow.WorkflowRef -> IO (Maybe Instance.WorkflowStatus)+workflowStatus store ref = do+ result <- runStoreIO store (Instance.lookupInstance (WorkflowName ref.workflowName) (WorkflowId ref.workflowId))+ either (fail . show) (pure . fmap (.status)) result++workflowLeaseOwner :: KirokuStore -> OpsWorkflow.WorkflowRef -> IO (Maybe Text)+workflowLeaseOwner store ref = do+ result <- runStoreIO store (Instance.lookupInstance (WorkflowName ref.workflowName) (WorkflowId ref.workflowId))+ either (fail . show) (pure . (>>= (.leasedBy))) result++awakeableStatus :: KirokuStore -> UUID.UUID -> IO (Maybe Awakeable.AwakeableStatus)+awakeableStatus store awakeableId = do+ result <- runStoreIO store (Awakeable.lookupAwakeable awakeableId)+ either (fail . show) (pure . fmap (.status)) result++timerStatus :: KirokuStore -> Timer.TimerId -> IO (Maybe Timer.TimerStatus)+timerStatus store timerId = do+ result <- runStoreIO store (Timer.lookupTimer timerId)+ either (fail . show) (pure . fmap (.status)) result++timerRequest :: String -> UTCTime -> Timer.TimerRequest+timerRequest rawId fireAt =+ Timer.TimerRequest+ { timerId = Timer.TimerId (maybe (error "test timer UUID") Function.id (UUID.fromString rawId)),+ processManagerName = "billing",+ correlationId = "invoice-1",+ fireAt,+ payload = object ["kind" .= ("reminder" :: Text)]+ }++resultArrayLength :: OpsOutcome -> Maybe Int+resultArrayLength = \case+ Succeeded OpsResult {jsonValue = Aeson.Array values} -> Just (Vector.length values)+ _ -> Nothing++resultArrayLengthFrom :: Key -> OpsOutcome -> Maybe Int+resultArrayLengthFrom key = \case+ PreviewRequired OpsResult {jsonValue = Aeson.Object value} _ ->+ case KeyMap.lookup key value of+ Just (Aeson.Array values) -> Just (Vector.length values)+ _ -> Nothing+ _ -> Nothing++resultArrayLengthFromObject :: Key -> OpsOutcome -> Maybe Int+resultArrayLengthFromObject key = \case+ Succeeded OpsResult {jsonValue = Aeson.Object value} ->+ case KeyMap.lookup key value of+ Just (Aeson.Array values) -> Just (Vector.length values)+ _ -> Nothing+ _ -> Nothing++resultCount :: OpsOutcome -> Maybe Int+resultCount = fmap fromIntegral . jsonInteger "count"++jsonInteger :: Key -> OpsOutcome -> Maybe Int64+jsonInteger key = \case+ Succeeded OpsResult {jsonValue = Aeson.Object value} -> numberAt key value+ _ -> Nothing++humanField :: Text -> OpsOutcome -> Maybe Text+humanField key = \case+ Succeeded OpsResult {headers, rows = [row]} -> lookup key (zip headers row)+ _ -> Nothing++jsonStringArray :: Key -> OpsOutcome -> Maybe [Text]+jsonStringArray key = \case+ Succeeded OpsResult {jsonValue = Aeson.Object value} -> do+ Aeson.Array values <- KeyMap.lookup key value+ traverse+ ( \case+ Aeson.String item -> Just item+ _ -> Nothing+ )+ (Vector.toList values)+ _ -> Nothing++jsonIntegerFromPreview :: Key -> OpsOutcome -> Maybe Int64+jsonIntegerFromPreview key = \case+ PreviewRequired OpsResult {jsonValue = Aeson.Object value} _ -> numberAt key value+ _ -> Nothing++numberAt :: Key -> KeyMap.KeyMap Aeson.Value -> Maybe Int64+numberAt key value = do+ Aeson.Number number <- KeyMap.lookup key value+ pure (floor number)++jsonBool :: Key -> OpsOutcome -> Maybe Bool+jsonBool key = \case+ Succeeded OpsResult {jsonValue = Aeson.Object value} -> do+ Aeson.Bool result <- KeyMap.lookup key value+ pure result+ _ -> Nothing++firstWorkflowId :: OpsOutcome -> Maybe Text+firstWorkflowId = \case+ Succeeded OpsResult {jsonValue = Aeson.Array values} -> do+ Aeson.Object first <- values Vector.!? 0+ Aeson.String workflowId <- KeyMap.lookup "workflow_id" first+ pure workflowId+ _ -> Nothing++journalEventCount :: OpsOutcome -> Maybe Int+journalEventCount = \case+ Succeeded OpsResult {jsonValue = Aeson.Object value} ->+ case KeyMap.lookup "events" value of+ Just (Aeson.Array events) -> Just (Vector.length events)+ _ -> Nothing+ _ -> Nothing++isPreview :: OpsOutcome -> Bool+isPreview PreviewRequired {} = True+isPreview _ = False++isSucceeded :: OpsOutcome -> Bool+isSucceeded Succeeded {} = True+isSucceeded _ = False++isLeft :: Either a b -> Bool+isLeft Left {} = True+isLeft Right {} = False++keiroOpsExecutable :: IO FilePath+keiroOpsExecutable = do+ (exitCode, stdoutText, stderrText) <-+ readProcessWithExitCode "cabal" ["list-bin", "exe:keiro-ops"] ""+ case exitCode of+ ExitSuccess -> pure (Text.unpack (Text.strip (Text.pack stdoutText)))+ ExitFailure code -> fail ("cabal list-bin keiro-ops failed (" <> show code <> "): " <> stderrText)++executeSql :: Text -> Text -> IO ()+executeSql connectionString sql =+ bracket acquire Hasql.release $ \connection -> do+ result <- Hasql.use connection (HasqlSession.script sql)+ either (fail . show) pure result+ where+ acquire = do+ result <- Hasql.acquire (HasqlSettings.connectionString connectionString)+ either (fail . show) pure result