salmon-ops-0.1.0.0: src/Salmon/Actions/Serve.hs
{-# LANGUAGE OverloadedStrings #-}
{-# LANGUAGE ScopedTypeVariables #-}
{- | A long-running convergence loop, fed a stream of seeds.
Where @run up@ is a one-shot "expand this one directive and walk it once",
'serve' keeps a 'World' around: a set of seeds that have been declared, the
graphs those seeds evaluated to, and — unified across all of them by 'Ref' —
a per-node 'NodeState' saying which 'Direction' that node is wanted in and
whether it has 'Converged' there yet.
The unit of input is a /declaration/: a seed, plus what to do with it (see
'ServeCommand'). Declaring a seed evaluates it and folds the result into two
structures: 'worldMagma', one representative per node keyed by 'Ref' (what
each node /is/), and 'worldLedger', a "Salmon.Op.Ledger" entry per
declaration saying which nodes and which precedence edges that declaration
asks for. A @down@ retracts its entry rather than deleting it, because a
retracted declaration's /edges/ are exactly what says what order to take its
nodes down in. From the ledger everything else is derived:
* a node some live declaration asks for is wanted 'TurnUp';
* a node this world still tracks and no live declaration still asks for is
wanted 'TurnDown';
* flipping a node's direction resets it to 'Pending', so it gets applied
again in the new direction.
An 'Epoch' — the seed, the directive, and the graph it evaluated to — is
kept only while its declaration is live, because the only thing still needing
a graph is the up pass. The teardown works off the magma and the ledger, so
storage is bounded by node count and by how many declarations are live or
retiring, not by the shape of what has been declared.
Convergence then runs — after every declaration, and on demand via
@converge@ — as one teardown pass followed by one bring-up pass, each of
which is just "Salmon.Actions.UpDown".'UpDown.downTreeWith' /
'UpDown.upTreeWith' over the relevant graphs with a 'UpDown.Gate' that
filters down to the nodes wanted in that pass and not yet converged. The
dependency ordering, the dedup-by-'Ref', and the "a failed node blocks
whatever depended on it" containment therefore behave exactly as they do for
@run up@ / @run down@; the only thing this module adds on top is the memory
of what has already been done. Nodes that end a pass 'Errored' or 'Blocked'
stay non-converged and are retried by the next pass.
That memory is bounded, which matters for a process meant to stay up: a
node leaves when it has settled down, a contribution leaves once none of its
nodes is still on its way down, and an epoch leaves as soon as its
declaration is retired or superseded. 'resettle' collects after every
declaration and every convergence; see 'prune' for the rules. What survives collection is
'worldLog', one small line per declaration ever made, which is what
@history@ prints — so the record of /what was declared/ outlives the graphs
that were declared.
== Between passes, the nodes are tended
A convergence pass is one attempt at whatever is outstanding, and then it is
over. That leaves a gap this loop used to have no answer for: an effect that
goes away on its own — a service that dies, a file something else deletes —
is not noticed until somebody types @converge@.
So while this loop is idle, every node it knows has a machine of its own
("Salmon.Actions.Upkeep"), running in whatever direction the node is wanted.
A node wanted up rechecks its own effect on an adaptive delay and runs @up@
again if it has gone; a node wanted down retries its @down@ until it works,
then stops.
/Idle/ is exactly the condition, and it is 'loop' that enforces it: the
machines start when nothing is waiting in the input and stand down before any
command is handled (stopping /waits for/ an @up@ or @down@ in flight rather
than interrupting one, since a half-applied effect is worse than a slow
command). Two things fall out, and the second is the reason:
* a piped script has every line, end-of-input included, queued before the
first pass finishes, so it is never supervised at all — @serve \< script@
stays a deterministic sequence of passes;
* there is nothing to race. Starting machines and then stopping them
because a command had been sitting in the queue all along would make
"was this node acted on?" depend on thread timing.
@supervise off@ stops them for good and leaves every effect exactly as it is:
@off@ is not a teardown. Note also that a restricted @converge --select@
scopes the /pass/, not the standing watch — a node the pass skipped is still
tended once the loop goes idle.
== Where the lines come from
The loop reads one inbox. What fills it is a list of 'Producer's, each on a
thread of its own, each pushing 'Line's tagged with the 'Origin' that typed
them; 'serveWith' is the one-producer case, standard input, and behaves
exactly as it did when the loop read a 'Handle' directly. The inbox is still
the loop's whole notion of idle — machines are tended while it is empty and
stand down before any command, whoever typed it — and a piped script is
still every line queued before the first pass ends. Two producers
interleave at line granularity and nothing more: a command is handled whole
before the next is read, and the order in which two producers' lines land
in the inbox is the order they are handled.
One decision is new with the list. /Only standard input's end of input ends
the loop./ Any other producer's 'Eof' is not a command — nothing is about to
act, so the machines are not stood down for it — and the loop reads on; a
socket client hanging up or a fetcher going quiet must not take the server
with it. A loop with no 'Stdin' producer at all therefore ends only on
@quit@. Such a producer hanging up is reported, as 'HungUp', at the moment
the loop reads its 'Eof' — which, the inbox being one queue, is after every
line it typed has been handled. Whoever holds a connection for that origin
can close it on that report and know nothing typed on it is still pending.
== Whose report is it
A report emitted while a command is handled belongs to whoever typed the
command; one emitted between commands (the tending machines' own) belongs
to nobody in particular. 'serveAttributed' says which, by stamping every
report with the 'Origin' of the line being handled — 'Nothing' outside a
command — through 'Attributed'. That is what lets a second client be
answered on its own connection rather than on the loop's standard output
("Salmon.Actions.Serve.Socket"); 'serveProducers' is the same loop with the
stamp thrown away.
This is what replaced @serveWakingWith@, an "these nodes want attention" hook
nothing in the repo ever drove. It was there because a node had no state of
its own to block on; now one does, so the hook is not a smaller version of
this — it is unnecessary. Restart policy and watchdogs are per-node and live
in "Salmon.Op.Supervision".
-}
module Salmon.Actions.Serve (
-- * Running
serve,
serveWith,
serveProducers,
serveAttributed,
serveFollowing,
serveObserved,
Followed (..),
AppliedDocument (..),
Mode (..),
renderMode,
-- * Input producers
Producer (..),
Line (..),
Origin (..),
Provenance (..),
renderOrigin,
originName,
handleProducer,
stdinProducer,
-- * Attributing reports
Attributed (..),
-- * Input language
ServeCommand (..),
Declaration (..),
Selection (..),
noSelection,
Topic,
parseServeCommand,
parseSelection,
tokenize,
-- * World state
World (..),
Collision (..),
emptyWorld,
Epoch (..),
EpochId (..),
LogEntry (..),
worldLogLimit,
NodeState (..),
Direction (..),
Convergence (..),
-- * Reading a world
worldDag,
worldPaths,
historyLinesMatching,
-- * Reporting
Report (..),
reportText,
renderReport,
) where
import Control.Comonad.Cofree (Cofree)
import Control.Concurrent (forkIO, killThread)
import Control.Concurrent.STM (TChan, TVar, atomically, isEmptyTChan, newTChanIO, readTChan, writeTChan)
import Control.Exception (IOException, SomeException, finally, try)
import Control.Monad (forM_, unless, when)
import Data.Aeson (FromJSON (..), ToJSON (..), Value, eitherDecode, encode, object, withObject, (.:), (.=))
import Data.Aeson.Types (parseEither)
import Data.ByteString.Lazy (ByteString)
import qualified Data.ByteString.Lazy as LByteString
import Data.Char (isSpace)
import Data.Foldable (traverse_)
import Data.Maybe (isJust)
import Data.IORef (IORef, atomicModifyIORef', modifyIORef', newIORef, readIORef, writeIORef)
import Data.List (nub, sortOn)
import qualified Data.List as List
import Data.Map.Strict (Map)
import qualified Data.Map.Strict as Map
import Data.Set (Set)
import qualified Data.Set as Set
import Data.Text (Text)
import qualified Data.Text as Text
import qualified Data.Text.IO as Text
import Data.Time.Clock (UTCTime)
import System.IO (Handle, hFlush, hGetLine, hIsEOF, stdout)
import qualified Salmon.Actions.Query as Query
import qualified Salmon.Actions.Concurrent as Concurrent
import qualified Salmon.Actions.UpDown as UpDown
import Salmon.Actions.UpDown (Requirement (..))
-- imported with their field selectors: OverloadedRecordDot only solves
-- HasField for fields that are in scope.
import Salmon.Builtin.Extension (Extension (..), Op, Track', evalDeps)
import Salmon.Op.Actions (Act (..), ShortHand)
import Salmon.Op.Concurrency (ConcurrencyLimit)
import Salmon.Op.Configure (Configure, gen)
import Salmon.Op.Graph (Graph)
import Salmon.Op.Dag (Dag)
import qualified Salmon.Op.Dag as Dag
import Salmon.Op.Ledger (Ledger)
import qualified Salmon.Op.Ledger as Ledger
import qualified Salmon.Op.Mailbox as Mailbox
import Salmon.Op.Ref (Ref, unRef)
import Salmon.Op.Rewrite (Phase (..), Rewrite, Rewritten)
import qualified Salmon.Op.Rewrite as Rewrite
-- 'Direction' used to be declared here, identically. It belongs with the
-- per-node state now that a node has state of its own, and is re-exported
-- from this module so nothing that named 'Serve.TurnUp' had to change.
import Salmon.Op.Status (Direction (..))
import qualified Salmon.Op.Status as MachineStatus
import Salmon.Op.Supervision (Micros (..))
import Salmon.Op.Track (run)
import Salmon.Reporter
import qualified Salmon.Actions.Upkeep as Upkeep
-------------------------------------------------------------------------------
-- | How far a node is from its wanted 'Direction'.
data Convergence
= -- | never applied in the current direction (new node, or the
-- direction just flipped under it)
Pending
| -- | (I6): was 'Converged', but a later declaration replaced this
-- 'Ref''s representative with one 'Dag.sameRepresentative' calls
-- different — so what it means to be converged may have changed too.
-- Treated exactly like 'Pending' by 'gateFor' (anything but
-- 'Converged' is worth a pass's attention): the node is handed to
-- 'UpDown.upDag' again, which asks its own @check@ before doing
-- anything, same as ever. A node with a content-comparing @check@
-- (e.g. 'Salmon.Builtin.Nodes.Filesystem.filecontents') settles
-- straight back to 'Converged' at the cost of one @check@ if the new
-- declaration didn't actually change what it writes; a node with no
-- @check@ gets exactly what it already gets under a bare 'Pending' —
-- an unconditional @up@, which the idempotency convention every node
-- author is already asked to follow makes safe. Kept as its own
-- constructor rather than folded into 'Pending' so @status@ can tell
-- "never touched" from "was up, now re-verifying".
Stale
| -- | applied in the current direction, or found to already be there
Converged
| -- | the last attempt threw; will be retried
Errored
| -- | the last attempt never ran because a neighbour failed; will be retried
Blocked
deriving (Show, Eq, Ord)
{- | What 'serve' remembers about one node of the unified graph. Keyed by
'Ref', so the same node reached through several seeds' graphs is one entry.
-}
data NodeState = NodeState
{ nodeShorthand :: !ShortHand
, nodeHelp :: !Text
, nodeDirection :: !Direction
, nodeConvergence :: !Convergence
, nodeStatus :: !(Maybe MachineStatus.Status)
-- ^ (R3). What this node's machine last had to say for itself: its own
-- 'Salmon.Actions.UpDown.CheckResult', how long ago it last did
-- anything observable, and its ring of output — the same 'Status' a
-- machine's neighbours block on while tending is running, snapshotted
-- by 'stopTending' at the one moment it is readable from outside: after
-- the machine has stood down (or been detached into
-- 'Salmon.Actions.Upkeep.Kept') but before the 'Upkeep.Supervisor'
-- holding its 'TVar' is dropped. 'Nothing' for a node that has never
-- been tended — declared while supervision is off, or not yet reached
-- by a first idle pass.
--
-- Freshness rides the existing rhythm rather than adding one: every
-- command runs 'stopTending' first (see 'loop'), so a node that /was/
-- tended has a snapshot from mere moments before whatever just read it.
-- A holding machine is adopted back into the next 'Upkeep.Supervisor'
-- the next time tending starts, which is what keeps its snapshot
-- refreshing across commands too, rather than freezing at whenever it
-- first started holding.
}
deriving (Show)
newtype EpochId = EpochId {unEpochId :: Int}
deriving (Show, Eq, Ord)
{- | One declaration, and everything derived from it at the time it was made.
An epoch is kept only while its declaration is /live/, because the only thing
it is still needed for is @--select@ resolution, which is about active seeds
and needs the graph's paths rather than the magma's nodes. What a retired declaration leaves
behind is its 'Ledger.Contribution' — two flat sets — plus its nodes in
'worldMagma', which is all a teardown needs and is bounded by node count
rather than by graph shape. Re-declaring the same seed appends a new epoch
and 'prune' drops the superseded one; 'worldLog' remembers that it happened.
-}
data Epoch seed directive = Epoch
{ epochId :: !EpochId
, epochDeclaration :: !Declaration
, epochDirection :: !Direction
, -- | who made this declaration: typed, loaded from a file, or fetched
-- from a registry. Kept for @history@.
epochOrigin :: !Origin
, -- | the argv this seed was declared with, kept for @history@; a
-- directive-file declaration (@up-directive@ and friends) gets a
-- synthetic @["<directive-file>", path]@ here instead.
epochTokens :: [String]
, -- | 'Nothing' when this epoch was declared straight from a directive
-- file, which has no seed to keep.
epochSeed :: Maybe seed
, epochDirective :: directive
, -- | identity of the seed for the active set: its encoded directive, so
-- that two spellings of the same desired state are one active seed
epochKey :: !ByteString
, -- | the graph as evaluated when the seed was declared. The one thing
-- left that needs a graph rather than the magma: @--select@ resolves
-- path patterns, and a 'Dag' has 'Ref's and edges but no paths.
epochGraph :: Cofree Graph Op
}
{- | What @history@ prints, and all that is kept of an 'Epoch' once 'prune'
has collected it. Deliberately small — no graph, no seed, no directive —
which is what makes it affordable to keep long after the epoch itself is
gone. 'logRefs' is the one non-trivial field: the node set the declaration
contributed, kept so @history --select@ still answers for a collected epoch.
-}
data LogEntry = LogEntry
{ logEpoch :: !EpochId
, logDeclaration :: !Declaration
, logOrigin :: !Origin
, logTokens :: [String]
, logRefs :: !(Set Ref)
}
deriving (Show)
-- | How many declarations 'worldLog' remembers.
worldLogLimit :: Int
worldLogLimit = 1000
data World seed directive = World
{ worldNextId :: !Int
, -- | the epochs whose graphs a future pass could still walk — the
-- active ones, plus retired ones that still describe a node to turn
-- down. Newest first; collected by 'prune'.
worldEpochs :: [Epoch seed directive]
, -- | one line per declaration ever made, newest first, outliving the
-- epoch's graph. Capped at 'worldLogLimit'.
worldLog :: [LogEntry]
, -- | declarations that have fallen off the end of 'worldLog'
worldLogDropped :: !Int
, -- | which declarations still want which nodes and edges — live ones
-- (what should be up) and retiring ones (whose edges are still the
-- only description of what order to take their nodes down in). Keyed
-- by 'epochKey', so two spellings of one desired state are one entry.
-- This is what replaced keeping a retired seed's whole graph.
worldLedger :: !(Ledger ByteString)
, -- | one representative per node, merged across every declaration that
-- has mentioned it (last writer wins, see "Salmon.Op.Dag"). Holds what
-- a node /is/ — its @up@\/@check@\/@down@ — where 'worldNodes' holds
-- where it has got to. Never an 'Op': that would retain the whole
-- expanded closure and bound nothing.
worldMagma :: !(Map Ref (Act Extension))
, -- | the nodes whose current representative won a collision that is
-- still standing: what @\/dag@ shows beside such a node so a client
-- that did not catch the pass's 'UpDown.Conflicting' can still show
-- the pair. See 'Collision' for when an entry appears and goes.
worldConflicts :: !(Map Ref Collision)
, -- | every node still being managed or still to be torn down, unified
-- by 'Ref'. A node that has converged 'TurnDown' is finished and is
-- dropped, so a fully retired world settles empty.
worldNodes :: Map Ref NodeState
}
{- | The machines currently tending this world's nodes, and whether they are
wanted at all.
Deliberately not part of 'World': a 'World' is a pure value this module hands
back to its caller, and a running 'Upkeep.Supervisor' is neither pure nor
meaningful once 'serve' has returned.
-}
data Tending = Tending
{ tendingSup :: !(IORef (Maybe (Upkeep.Supervisor Extension)))
, tendingKept :: !(IORef (Upkeep.Kept Extension))
-- ^ machines still holding a 'Salmon.Builtin.Extension.managed' effect
-- up, between one supervisor and the next. These outlive a command
-- precisely because stopping a supervisor means "stop tending", and a
-- @status@ that killed every service would be a poor reading of that.
-- See 'Salmon.Actions.Upkeep.Kept'.
, tendingOn :: !(IORef Bool)
-- ^ @supervise off@ clears this; nothing is tended between passes, and
-- @serve@ behaves as it did before per-node machines existed.
, tendingAutoConverge :: !(IORef Bool)
-- ^ @autoconverge off@ clears this: a declaring command
-- (@up@\/@only@\/@down@\/@clear@\/the @-directive@ forms) still records
-- the epoch and updates 'worldNodes'/'worldLedger' as usual, but the
-- convergence pass that would otherwise follow it immediately is
-- skipped, leaving whatever @status@\/@query@ already show unchanged
-- until an explicit @converge@. On by default, matching every existing
-- caller's behaviour.
, tendingPending :: !(IORef (Map Ref [Mailbox.Instruction]))
-- ^ (R2). instructions an operator posted while no supervisor was
-- running to hand them to. A one-shot machine does not survive a
-- command the way a holding one does (see 'tendingKept'), so
-- @force@\/@recheck@\/@pause@\/@resume@ cannot post straight into a
-- mailbox that is about to be discarded — 'startTending' delivers these
-- the moment the /next/ supervisor's machines exist (both freshly
-- started and adopted), then clears the queue. "Force this node next
-- time you look at it" rather than keeping every one-shot machine alive
-- just so it has a mailbox to post into.
}
{- | A representative that lost to the magma's current one, kept for as
long as somebody still wants the loser's version.
Two kinds of collision land here, through 'record'. One is inside a single
declaration: the graph reaches one 'Ref' from two differently-described
nodes, which is 'Dag.dagConflicts' and has always been reported
'UpDown.Conflicting' at declare time. The other is /across/ declarations:
this declaration describes a 'Ref' differently from what the magma holds,
and another live declaration still wants that 'Ref' — two seeds colliding
on one node, which last-writer-wins resolves silently otherwise (the same
key re-declared with a change is not a collision but (I6)'s 'Stale').
'collisionHolders' is who was standing on the losing side — the other live
declarations for the cross kind, the declaration itself for the inside kind
— and is what keeps the entry honest without keeping every declaration's
representative around: a re-declaration that leaves the magma's
representative as it is keeps the entry while a holder is still live, a
re-declaration that changes it recomputes, and 'prune' drops an entry whose
holders have all retired or whose node has left the magma.
-}
data Collision = Collision
{ collisionConflict :: !(Dag.Conflict Extension)
-- ^ 'Dag.conflictKept' is the magma's representative at the time of
-- the write, 'Dag.conflictReplaced' the one it beat
, collisionHolders :: !(Set ByteString)
-- ^ 'epochKey's of the declarations on the losing side
}
emptyWorld :: World seed directive
emptyWorld = World 0 [] [] 0 Ledger.emptyLedger Map.empty Map.empty Map.empty
-------------------------------------------------------------------------------
-- | What declaring a seed does to the active set.
data Declaration
= -- | @up@: add this seed to the active set
Add
| -- | @only@: make this seed the whole active set, retiring the others
Replace
| -- | @down@: retire this seed
Remove
deriving (Show, Eq, Ord)
{- | A pair of select\/exclude path-patterns, understood the same way as
"Salmon.Actions.Query" ('Query.parsePattern'\/'Query.matchPattern'). An empty
'selSelect' means "everything" (mirrors 'Query.resolveSelectors'); an empty
'selExclude' subtracts nothing. 'noSelection' is both empty, and is what a
bare @status@\/@history@\/@converge@\/@query@ (no @--select@\/@--exclude@ at
all) parses to.
-}
data Selection = Selection
{ selSelect :: ![Text]
, selExclude :: ![Text]
}
deriving (Show, Eq)
noSelection :: Selection
noSelection = Selection [] []
data ServeCommand
= Declare !Declaration ![String]
| -- | @up-directive@\/@only-directive@\/@down-directive@: declare a seed
-- straight from a directive JSON file, skipping seed-arg parsing.
DeclareDirective !Declaration !FilePath
| -- | the same declaration with the directive's JSON already in hand
-- rather than in a file. Not spelled by any line of the input language
-- ('parseServeCommand' never produces it); it exists for a producer
-- that holds a document with a directive in it ("Salmon.Actions.Follow")
-- and would otherwise have to write that directive to a file to name
-- it. The 'Text' is what @history@ prints in place of an argv.
DeclareInline !Declaration !Text !Value
| -- | @load@: run a file of serve-command lines, in order, as if typed.
Load !FilePath
| -- | @clear@: retire every seed (everything known goes down)
Clear
| -- | @converge@: re-attempt whatever has not converged; a non-empty
-- 'Selection' restricts this one pass to matching nodes only.
Converge !Selection
| Status !Selection
| History !Selection
| -- | @query@: annotate the world's nodes with a 'Selection', without
-- acting on anything.
QueryCmd !Selection
| -- | @supervise on@\/@supervise off@: whether to keep tending nodes
-- between convergence passes. On by default.
Supervise !Bool
| -- | @autoconverge on@\/@autoconverge off@: whether a declaring
-- command (@up@\/@only@\/@down@\/@clear@\/the @-directive@ forms)
-- triggers a convergence pass on its own. On by default; @off@ lets
-- several declarations (or an inspection via @status@\/@query@) sit
-- between the declaration and an explicit @converge@.
AutoConverge !Bool
| -- | @force@\/@recheck@\/@pause@\/@resume@ [--select P]... [--exclude
-- P]...: queue a 'Mailbox.Instruction' for the matching nodes, to be
-- delivered the next time this world's nodes are tended (R2). An
-- empty selection means every node, same as @status@\/@query@.
Instruct !Mailbox.Instruction !Selection
| -- | @fetch@: ask the fetcher ("Salmon.Actions.Follow") for a round
-- now — its ladder forgotten, whatever it has pending injected as soon
-- as the round is over — rather than at its next scheduled one. The
-- loop cannot call into a producer, so it pulls a hook 'serveFollowing'
-- was given; without one, nothing is being followed and it says so.
Fetch
| -- | @help@\/@help TOPIC@: print the command reference, or (when
-- 'Just' a recognised 'Topic') a lengthier explanation of just that
-- one command. 'Nothing', or a topic 'lookupTopic' doesn't recognise,
-- both fall back to the same full reference.
Help !(Maybe Topic)
| Quit
| -- | blank line or comment
Noop
deriving (Show, Eq)
-- | A @help@ argument, matched case-insensitively against 'helpTopics'.
type Topic = Text
declarationDirection :: Declaration -> Direction
declarationDirection Add = TurnUp
declarationDirection Replace = TurnUp
declarationDirection Remove = TurnDown
{- | Parse one line of 'serve' input: a command word followed, for the
declaring commands, by the seed's own command-line arguments.
-}
parseServeCommand :: String -> Either Text ServeCommand
parseServeCommand line =
case dropWhile isSpace line of
[] -> Right Noop
('#' : _) -> Right Noop
_ -> dispatch =<< tokenize line
where
dispatch toks =
case toks of
[] -> Right Noop
(w : args) ->
case w of
"up" -> Right (Declare Add args)
"only" -> Right (Declare Replace args)
"down" -> Right (Declare Remove args)
"up-directive" -> onlyFile w args (DeclareDirective Add)
"only-directive" -> onlyFile w args (DeclareDirective Replace)
"down-directive" -> onlyFile w args (DeclareDirective Remove)
"load" -> onlyFile w args Load
"clear" -> nullary w args Clear
"converge" -> Converge <$> parseSelection args
"status" -> Status <$> parseSelection args
"history" -> History <$> parseSelection args
"query" -> QueryCmd <$> parseSelection args
"supervise" -> onOff w args Supervise
"autoconverge" -> onOff w args AutoConverge
"force" -> Instruct Mailbox.Force <$> parseSelection args
"recheck" -> Instruct Mailbox.Recheck <$> parseSelection args
"pause" -> Instruct Mailbox.Pause <$> parseSelection args
"resume" -> Instruct Mailbox.Resume <$> parseSelection args
"fetch" -> nullary w args Fetch
"help" -> Help <$> helpTopic w args
"?" -> Help <$> helpTopic w args
"quit" -> nullary w args Quit
"exit" -> nullary w args Quit
_ -> Left ("unknown command: " <> Text.pack w)
nullary w args cmd
| null args = Right cmd
| otherwise = Left (Text.pack w <> " takes no argument")
helpTopic w args =
case args of
[] -> Right Nothing
[t] -> Right (Just (Text.pack t))
_ -> Left (Text.pack w <> " takes at most one topic argument")
onlyFile w args mk =
case args of
[path] -> Right (mk path)
_ -> Left (Text.pack w <> " takes exactly one file argument")
onOff w args mk =
case args of
["on"] -> Right (mk True)
["off"] -> Right (mk False)
_ -> Left (Text.pack w <> " takes exactly one of `on` or `off`")
{- | Scans a token list for repeated @--select PATTERN@\/@--exclude
PATTERN@ pairs, shared by @query@\/@status@\/@history@\/@converge@.
-}
parseSelection :: [String] -> Either Text Selection
parseSelection = go [] []
where
go sel exc [] = Right (Selection (reverse sel) (reverse exc))
go sel exc ("--select" : p : rest) = go (Text.pack p : sel) exc rest
go sel exc ("--exclude" : p : rest) = go sel (Text.pack p : exc) rest
go _ _ ["--select"] = Left "--select needs a PATTERN argument"
go _ _ ["--exclude"] = Left "--exclude needs a PATTERN argument"
go _ _ (w : _) = Left ("unrecognized argument: " <> Text.pack w)
{- | Split a line into argv-style tokens, honouring single quotes, double
quotes and backslash escapes, so a seed can carry values with spaces in them.
-}
tokenize :: String -> Either Text [String]
tokenize = outside []
where
outside toks s =
case s of
[] -> Right (reverse toks)
(c : cs)
| isSpace c -> outside toks cs
| otherwise -> word toks "" (c : cs)
word toks cur s =
case s of
[] -> Right (reverse (reverse cur : toks))
(c : cs)
| isSpace c -> outside (reverse cur : toks) cs
| c == '\\' -> escape (word toks) cur cs
| c == '\'' -> quoted '\'' toks cur cs
| c == '"' -> quoted '"' toks cur cs
| otherwise -> word toks (c : cur) cs
quoted q toks cur s =
case s of
[] -> Left "unterminated quote"
(c : cs)
| c == q -> word toks cur cs
| c == '\\' && q == '"' -> escape (quoted q toks) cur cs
| otherwise -> quoted q toks (c : cur) cs
escape k cur s =
case s of
(d : ds) -> k (d : cur) ds
[] -> Left "trailing backslash"
-------------------------------------------------------------------------------
data Report
= Started
| -- | input closed
Stopped
| -- | a producer other than standard input has no more lines, and every
-- line it did have has been handled. Never for 'Stdin', whose end of
-- input is 'Stopped'.
HungUp !Origin
| BadCommand !Text
| BadSeed !Text
| BadDirective !Text
| -- | reading a @load@ file failed, or its nesting was too deep
BadLoad !Text
| Loading !FilePath
| -- | lines run from a @load@ file
LoadDone !FilePath !Int
| -- | epoch, direction, nodes in its graph, live declarations afterwards
Declared !EpochId !Direction !Int !Int
| -- | number of seeds retired
Cleared !Int
| -- | @supervise on@\/@supervise off@
Supervised !Bool
| -- | @autoconverge on@\/@autoconverge off@
AutoConverged !Bool
| -- | (R2). a @force@\/@recheck@\/@pause@\/@resume@ was queued for this
-- many nodes; takes effect once tending next starts, not immediately
Instructed !Mailbox.Instruction !Int
| -- | @fetch@: whether anything is being followed (a round was asked
-- for), or not (nothing to ask)
FetchRequested !Bool
| -- | something a node's own machine had to say between convergence
-- passes. See 'Salmon.Actions.Upkeep.Report'; the chatty half of that
-- stream is filtered out before it reaches here.
Tended !(Upkeep.Report Extension)
| -- | nodes to turn down, nodes to turn up
ConvergeStart !Int !Int
| -- | everything applied cleanly, nodes still not converged
ConvergeStop !Bool !Int
| -- | the loop's 'Mode', then nodes, plus every live declaration's
-- path(s) to each one (see 'worldPaths') — the thing a
-- @--select@\/@--exclude@ pattern is actually built from.
StatusReport !Mode ![(Ref, NodeState)] !(Map Ref [Text])
| -- | epoch, declaration, still active, who declared it, argv
HistoryReport ![(EpochId, Declaration, Bool, Origin, [String])]
| -- | declarations too old to still be in 'worldLog'; emitted after a
-- 'HistoryReport' so @history@ never silently claims to be complete
HistoryElided !Int
| -- | world nodes annotated against a resolved selection: selected, excluded, paths
QueryReport ![(Ref, NodeState)] !(Set Ref) !(Set Ref) !(Map Ref [Text])
| -- | @help@: the full reference ('Nothing', or a 'Topic' 'lookupTopic'
-- didn't recognise), or a lengthier explanation of just that one
-- recognised 'Topic'.
HelpText !(Maybe Topic)
| -- | the status sink ("Salmon.Actions.Serve.StatusSink") could not
-- write its document: path, why. Emitted from the sink's own thread,
-- once per run of failures rather than once per attempt, and never
-- attributed to a typist; the loop keeps serving.
SinkFailed !FilePath !Text
deriving (Show)
-- | Prints 'Report's in a human-readable, one-event-per-block form.
reportText :: Reporter Report
reportText = ReporterM $ \rep -> do
-- one write per report rather than one per line: another producer's
-- reporter ("Salmon.Actions.Follow") shares this handle from its own
-- thread, and two half-lines interleaved are not two reports.
Text.putStr (Text.unlines (renderReport rep))
hFlush stdout
renderReport :: Report -> [Text]
renderReport rep =
case rep of
Started ->
[ "serve: ready"
, "serve: type `help` for the command reference"
]
Stopped -> ["serve: input closed"]
HungUp origin -> ["serve: " <> originName origin <> " hung up"]
BadCommand err -> ["serve: " <> err]
BadSeed err -> ("serve: cannot configure seed:") : Text.lines err
BadDirective err -> ("serve: cannot decode directive:") : Text.lines err
BadLoad err -> ["serve: " <> err]
Loading path -> ["serve: loading " <> Text.pack path]
LoadDone path n -> ["serve: loaded " <> Text.pack path <> " (" <> tshow n <> " line(s))"]
Declared eid dir nnodes nactive ->
[ Text.unwords
[ "serve: epoch"
, renderEpochId eid
, renderDirection dir
, "(" <> tshow nnodes <> " nodes,"
, tshow nactive <> " active seed(s))"
]
]
Cleared n -> ["serve: retired " <> tshow n <> " seed(s)"]
Supervised True -> ["serve: supervising (nodes are tended between passes)"]
Supervised False -> ["serve: not supervising (nodes are left alone between passes)"]
AutoConverged True -> ["serve: auto-converging (each declaration converges immediately)"]
AutoConverged False -> ["serve: not auto-converging (declarations wait for an explicit `converge`)"]
FetchRequested True -> ["serve: fetching now"]
FetchRequested False -> ["serve: nothing is being followed (start with --follow to fetch declarations)"]
Instructed instr n ->
[ Text.unwords
[ "serve: queued"
, Text.toLower (tshow instr)
, "for"
, tshow n
, "node(s), to take effect once tending next starts"
]
]
Tended t -> renderTended t
ConvergeStart ndown nup ->
["serve: converging (" <> tshow ndown <> " down, " <> tshow nup <> " up)"]
ConvergeStop ok remaining ->
[ Text.unwords
[ "serve:"
, -- 'remaining' rather than 'ok' decides the headline: a
-- restricted pass can leave nodes pending (skipped, not
-- attempted) while still reporting 'ok' — nothing it
-- actually attempted failed.
if remaining == 0 then "converged" else "converge incomplete"
, "(" <> tshow remaining <> " node(s) left"
, if ok then ")" else ", including a failure)"
]
]
StatusReport mode [] _ -> ["serve: mode: " <> renderMode mode, "serve: no nodes"]
StatusReport mode xs paths -> ("serve: mode: " <> renderMode mode) : "serve: nodes:" : concatMap (renderNode paths) (sortOn statusOrder xs)
HistoryReport [] -> ["serve: no seed declared yet"]
HistoryReport xs -> "serve: seeds:" : fmap renderEpochLine xs
HistoryElided n -> ["serve: " <> tshow n <> " earlier declaration(s) elided"]
QueryReport [] _ _ _ -> ["serve: no nodes"]
QueryReport xs sel exc paths -> "serve: nodes:" : concatMap (renderQueryNode paths sel exc) (sortOn statusOrder xs)
HelpText mtopic ->
case mtopic >>= lookupTopic of
Just detailed -> detailed
Nothing -> commandReference
SinkFailed path err -> ("serve: status sink " <> Text.pack path <> " could not be written:") : Text.lines err
where
statusOrder :: (Ref, NodeState) -> (Direction, Convergence, ShortHand, Text)
statusOrder (r, st) = (st.nodeDirection, st.nodeConvergence, st.nodeShorthand, unRef r)
{- | One summary line, plus (R3) a trailing detail block for a node whose
last known check was 'Salmon.Actions.UpDown.Failure' — its own ring of
output, tail-capped so one wedged node cannot bury the rest of the
listing. A node this has never tended (never supervised, or not yet
reached by an idle pass) says so rather than showing stale silence as if
it meant something. -}
renderNode :: Map Ref [Text] -> (Ref, NodeState) -> [Text]
renderNode paths (r, st) = summary : pathLines ++ detail
where
summary =
Text.unwords
[ " "
, renderDirection st.nodeDirection
, Text.justifyLeft 9 ' ' (tshow st.nodeConvergence)
, Text.justifyLeft 10 ' ' (unRef r)
, st.nodeShorthand
, renderVerdict st.nodeStatus
]
pathLines = case Map.findWithDefault [] r paths of
[] -> [" path: (none — not reached by any live declaration's graph)"]
[p] -> [" path: " <> p]
ps -> " paths:" : [ " " <> p | p <- ps ]
detail = case st.nodeStatus of
Just ms | UpDown.Failure _ <- ms.statusCheck -> renderRingTail ms.statusOutput
_ -> []
renderVerdict :: Maybe MachineStatus.Status -> Text
renderVerdict Nothing = "[not yet tended]"
renderVerdict (Just ms) = "[" <> tshow ms.statusCheck <> "]"
-- | The last few lines of a node's output ring, oldest of the shown
-- ones first — enough to see what a failing node was last saying
-- without dumping the whole (up to 256-line) ring into a status listing.
renderRingTail :: MachineStatus.Ring -> [Text]
renderRingTail ring =
case MachineStatus.ringLines ring of
[] -> []
ls ->
let shown = drop (max 0 (length ls - ringTailLines)) ls
omitted = length ls - length shown
header
| omitted > 0 = " last output (" <> tshow omitted <> " earlier line(s) omitted):"
| otherwise = " last output:"
in header : fmap (" " <>) shown
ringTailLines :: Int
ringTailLines = 10
renderQueryNode :: Map Ref [Text] -> Set Ref -> Set Ref -> (Ref, NodeState) -> [Text]
renderQueryNode paths sel exc entry@(r, _) = case renderNode paths entry of
[] -> []
(summary : rest) -> (summary <> annotation) : rest
where
annotation
| r `Set.member` exc = " [excluded]"
| r `Set.member` sel = " [selected]"
| otherwise = ""
-- a typed line renders exactly as it did before origins existed; any
-- other origin is a trailing annotation, so the argv stays where an
-- operator's eye already looks for it.
renderEpochLine :: (EpochId, Declaration, Bool, Origin, [String]) -> Text
renderEpochLine (eid, decl, active, origin, toks) =
Text.unwords $
[ " "
, renderEpochId eid
, Text.justifyLeft 8 ' ' (renderDeclaration decl)
, if active then "[active]" else "[retired]"
, Text.pack (unwords toks)
]
++ [ann | Just ann <- [renderOrigin origin]]
{- | How @history@ names where a declaration came from: 'Nothing' for a typed
line (the common case, and the one every existing transcript shows), a
bracketed annotation otherwise. A fetched declaration names its registry,
label, document id and digest, which is the whole point of recording it —
see "Salmon.Actions.Follow".
-}
renderOrigin :: Origin -> Maybe Text
renderOrigin origin =
case origin of
Stdin -> Nothing
Origin name -> Just ("[via " <> name <> "]")
Loaded path -> Just ("[loaded " <> Text.pack path <> "]")
Fetched prov ->
Just $
Text.concat
[ "[fetched "
, prov.provRegistry
, " label="
, prov.provLabel
, " id="
, prov.provDocument
, " sha256="
, Text.take 12 prov.provDigest
, "]"
]
{- | The supervision events worth an operator's attention, one line each.
Everything a machine says about its own progress — state transitions, the
next check's delay — is dropped upstream in 'serve''s own reporter rather
than rendered small here: a per-node line on every nap is a trace, not a
report.
-}
renderTended :: Upkeep.Report Extension -> [Text]
renderTended t =
case t of
Upkeep.Supervising nup ndown ->
["serve: tending " <> tshow nup <> " node(s) up, " <> tshow ndown <> " down"]
-- not rendered: the machines stand down before every command,
-- including a `help`, and a line saying so each time is noise. That
-- they came back is what the next 'Upkeep.Supervising' says.
Upkeep.Retired _ -> []
Upkeep.Wedged act (Micros us) ->
[ "serve: "
<> act.shorthand
<> " has been silent for "
<> tshow (us `div` 1000)
<> "ms, past its watchdog"
]
Upkeep.Holding n -> ["serve: " <> tshow n <> " node(s) still holding an effect up"]
Upkeep.Unwedged act -> ["serve: " <> act.shorthand <> " is moving again"]
-- not rendered: a live tail is for a client of `/events`; on a
-- terminal it would interleave every daemon's stdout with the reports.
Upkeep.Output _ _ -> []
Upkeep.GaveUp act n ->
[ "serve: "
<> act.shorthand
<> " gave up after "
<> tshow n
<> " consecutive failures; force or recheck it to try again"
]
-- a machine taken over from the previous supervisor: worth a line,
-- because the alternative (a restart) would have been visible and an
-- operator should be able to tell which happened.
Upkeep.Adopted act -> ["serve: " <> act.shorthand <> " kept running"]
Upkeep.Released act -> ["serve: " <> act.shorthand <> " let go"]
-- worth a line even though it is a normal consequence of a
-- declared policy: it is the one thing in the supervisor that
-- touches a node nobody asked about, so an operator seeing work
-- happen on a node they did not expect should be able to find out
-- why from the same stream.
Upkeep.Demoted act dep ->
[ "serve: "
<> act.shorthand
<> " sent back to wait: "
<> unRef dep
<> " stopped being up"
]
Upkeep.Paused act -> ["serve: " <> act.shorthand <> " paused (its effect is untouched)"]
Upkeep.Resumed act -> ["serve: " <> act.shorthand <> " resumed"]
Upkeep.Policy act _ ignored ->
[ "serve: "
<> act.shorthand
<> " declares "
<> tshow (1 + length ignored)
<> " supervision policies; the first is in force"
]
Upkeep.Escaped act e ->
("serve: " <> act.shorthand <> "'s own machine threw:") : Text.lines (tshow e)
-- filtered out before they get here; listed so a new constructor is
-- a compile error rather than a silent omission.
Upkeep.Acted _ -> []
Upkeep.Upkeep{} -> []
Upkeep.Downkeep{} -> []
Upkeep.NextLook{} -> []
-- the same kind of thing as 'NextLook', and filtered for the same
-- reason: it is a machine saying what it is waiting on, which is
-- most nodes most of the time. `status` is where to see it.
Upkeep.Parked{} -> []
-- also filtered, and for the same reason as 'Parked' — it is
-- announced on every sleep of a node that declared
-- `Op.Supervision.supReapply`, which for a busy directory tree could
-- be every few seconds. `status` is where to see whether a node is
-- being reapplied rather than polled.
Upkeep.Reapplying{} -> []
Upkeep.Untended{} -> []
renderDirection :: Direction -> Text
renderDirection TurnUp = "up"
renderDirection TurnDown = "down"
renderDeclaration :: Declaration -> Text
renderDeclaration Add = "up"
renderDeclaration Replace = "only"
renderDeclaration Remove = "down"
renderEpochId :: EpochId -> Text
renderEpochId eid = "#" <> tshow eid.unEpochId
tshow :: (Show a) => a -> Text
tshow = Text.pack . show
-- | The full command reference, printed by @help@\/@?@ with no topic, or
-- with a topic 'lookupTopic' doesn't recognise.
commandReference :: [Text]
commandReference =
[ "serve: commands:"
, " up <seed args...> add this seed to the active set"
, " only <seed args...> make this seed the whole active set, retiring the others"
, " down <seed args...> retire this seed"
, " up-directive <file> like `up`, but from a directive JSON file (no seed parsing)"
, " only-directive <file> like `only`, but from a directive JSON file"
, " down-directive <file> like `down`, but from a directive JSON file"
, " load <file> run a file of these command lines, in order, as if typed"
, " clear retire every seed (everything known goes down)"
, " converge [--select P]... [--exclude P]..."
, " re-attempt whatever has not converged;"
, " with --select/--exclude, restrict this one pass to matching nodes"
, " status [--select P]... [--exclude P]..."
, " list nodes and their direction/convergence"
, " history [--select P]... [--exclude P]..."
, " list past declarations"
, " query [--select P]... [--exclude P]..."
, " annotate nodes [selected]/[excluded], without acting on anything"
, " supervise on|off whether to keep tending nodes between passes (default on)"
, " autoconverge on|off whether a declaration converges immediately (default on)"
, " force [--select P]... [--exclude P]..."
, " run `up` on matching nodes even though their check says not to"
, " recheck [--select P]... [--exclude P]..."
, " look at matching nodes now, rather than at their next delay"
, " pause [--select P]... [--exclude P]..."
, " stop tending matching nodes, without touching their effect"
, " resume [--select P]... [--exclude P]..."
, " start tending matching nodes again"
, " fetch (--follow) fetch the followed documents now, not at the next round"
, " help, ? [TOPIC] print this reference, or (given a topic) more about just it"
, " quit, exit leave the loop, changing nothing on the way out"
, "serve: --select/--exclude patterns are /-separated node-path globs (* one segment, ** any depth);"
, " may repeat; omitting --select entirely means everything."
, "serve: `help TOPIC` for more, where TOPIC is one of:"
, " up, directive, load, clear, converge, status, history, query, select, supervise,"
, " autoconverge, force, fetch"
]
{- | @help TOPIC@'s lookup table, matched case-insensitively (several names
may share one block of text, e.g. @up@\/@only@\/@down@ all point at
'declareHelp'). A topic not listed here falls back to 'commandReference'
(see 'lookupTopic').
-}
helpTopics :: [(Topic, [Text])]
helpTopics =
[ ("up", declareHelp)
, ("only", declareHelp)
, ("down", declareHelp)
, ("directive", directiveHelp)
, ("up-directive", directiveHelp)
, ("only-directive", directiveHelp)
, ("down-directive", directiveHelp)
, ("load", loadHelp)
, ("clear", clearHelp)
, ("converge", convergeHelp)
, ("status", statusHelp)
, ("history", historyHelp)
, ("query", queryHelp)
, ("supervise", superviseHelp)
, ("watchdog", superviseHelp)
, ("autoconverge", autoConvergeHelp)
, ("force", instructHelp)
, ("recheck", instructHelp)
, ("pause", instructHelp)
, ("resume", instructHelp)
, ("fetch", fetchHelp)
, ("select", selectHelp)
, ("exclude", selectHelp)
, ("pattern", selectHelp)
]
lookupTopic :: Topic -> Maybe [Text]
lookupTopic t = lookup (Text.toLower t) helpTopics
declareHelp :: [Text]
declareHelp =
[ "serve: up / only / down <seed args...>"
, ""
, " up <seed args...> parses <seed args...> with the seed's own command-line parser (the"
, " same words that would follow `config` on the command line),"
, " configures it into a directive, and adds the resulting epoch to the"
, " active set."
, " only <seed args...> like `up`, but also retires every other currently active seed. A"
, " node shared with a retiring seed (e.g. an enclosing directory) is"
, " left alone if the new seed still wants it too."
, " down <seed args...> retires this seed. Its nodes go down unless another active seed"
, " still wants them."
, ""
, " A seed is identified by its encoded directive, not by its argv spelling: re-declaring an"
, " already-active, unchanged seed is a no-op (nothing pending, nothing re-run)."
, ""
, " Every declaration converges automatically right after being recorded (as if `converge`"
, " had been typed next); it is never itself scoped by --select/--exclude. `autoconverge off`"
, " turns this off, so several declarations can be recorded and inspected (`status`/`query`)"
, " before an explicit `converge` acts on any of them — see `help autoconverge`."
, ""
, " See also: `help directive` (declaring from a pre-generated directive file instead of"
, " seed args), `help load` (batch-declaring several seeds from a script file)."
]
directiveHelp :: [Text]
directiveHelp =
[ "serve: up-directive / only-directive / down-directive <file>"
, ""
, " Exactly like `up`/`only`/`down`, but the seed's own command-line parser is skipped"
, " entirely: <file> is read and JSON-decoded straight into the directive, e.g. the output"
, " of `my-salmon config <seed-args> > configs/foo.json` saved ahead of time."
, ""
, " Useful when a directive was already generated once (or came from somewhere other than"
, " this binary's own seed parser) and there is no seed value to reconstruct here — `status`"
, " still shows these nodes normally, and `history` records the file path in place of argv."
, ""
, " A malformed or unreadable file reports an error and declares nothing."
]
loadHelp :: [Text]
loadHelp =
[ "serve: load <file>"
, ""
, " Reads <file> and runs each of its lines through this exact same command language, in"
, " order, as if they had been typed (or piped) at the prompt one at a time — including"
, " further `load` lines, blank lines, and `#`-comments."
, ""
, " A `quit`/`exit` inside a loaded file ends the whole serve session, not just the load."
, ""
, " Nested loads are capped at a small depth to catch a file that (directly or indirectly)"
, " loads itself; exceeding it reports an error rather than looping forever."
, ""
, " This is the way to turn a directory of saved scripts (each a sequence of `up`/`only`/"
, " `down`/`up-directive`/... lines) into one `load configs/whatever.txt` declaration."
]
clearHelp :: [Text]
clearHelp =
[ "serve: clear"
, ""
, " Retires every currently active seed in one step (equivalent to a `down` for each). Every"
, " node no seed wants any more goes down on the convergence pass that follows automatically."
, " Takes no arguments."
]
convergeHelp :: [Text]
convergeHelp =
[ "serve: converge [--select PATTERN]... [--exclude PATTERN]..."
, ""
, " Re-attempts whatever has not yet converged: one teardown pass over nodes wanted down,"
, " then one bring-up pass over nodes wanted up. This runs automatically after every"
, " declaration; a bare `converge` is for retrying after fixing whatever made a node error"
, " out, or after a wait for some external condition."
, ""
, " With --select/--exclude, this one pass is additionally restricted to nodes matching the"
, " resolved selection (see `help select`) — anything outside it is left exactly as it was,"
, " neither attempted nor marked converged, so a later unrestricted `converge` still picks it"
, " up. Omitting both flags converges everything pending, as before."
, ""
, " Note the restriction scopes the *pass*, not the world: with supervision on (the"
, " default), a node this pass skipped is still tended once the loop goes idle, and may be"
, " acted on then. `supervise off` first if you want a pass to be the only thing that"
, " touches anything."
, ""
, " The report's headline (`converged` vs. `converge incomplete`) reflects how many nodes are"
, " still left afterwards, not just whether anything attempted this pass failed — a"
, " restricted pass can report no failure while still leaving excluded nodes pending."
]
statusHelp :: [Text]
statusHelp =
[ "serve: status [--select PATTERN]... [--exclude PATTERN]..."
, ""
, " First says which mode the loop is in — `interactive` (nothing followed: every declaration"
, " was typed or loaded), `following` (the world is what the registry last said), or `replay`"
, " (the registry could not be reached at startup and the world is a cached document: the"
, " last one applied before the restart, until a round in which every label answers)."
, " Then lists every node this world is still concerned with, unified by Ref across every seed"
, " that shares it, with its wanted direction (up/down) and convergence (Pending/Stale/"
, " Converged/Errored/Blocked). A node that has finished going down is dropped, so a world whose seeds"
, " have all been retired and converged lists nothing at all — `history` still shows they"
, " were declared."
, ""
, " With no flags at all, lists everything, exactly as before --select/--exclude existed"
, " (including nodes still on their way down). With --select/--exclude given, narrows the"
, " listing to the resolved selection (see `help select`) — this can, unlike the unfiltered"
, " form, only show nodes belonging to a currently active seed."
, ""
, " Each line also carries what the node's own machine last had to say for itself (its"
, " check, in brackets) once it has been tended at least once; `[not yet tended]` means"
, " supervision has not reached it yet. A node whose last word was a failure additionally"
, " shows the tail of its own output ring underneath — what it was doing right before it"
, " failed, which is otherwise nowhere to see."
, ""
, " Underneath each summary line is the path (or paths, if more than one live seed's graph"
, " reaches the same node) a --select/--exclude PATTERN would match to name it — the same"
, " slash-separated form `run tree`/`query` print, pasteable straight back in. This is the"
, " only place those paths are discoverable at all; a node with none listed belongs to no"
, " currently active seed (it is on its way down after being retired)."
]
historyHelp :: [Text]
historyHelp =
[ "serve: history [--select PATTERN]... [--exclude PATTERN]..."
, ""
, " Lists the declarations made, newest last, each tagged [active]/[retired] and showing the"
, " original argv (or, for a directive-file declaration, the file path). This is a log of"
, " what was asked for, kept long after the graph a declaration built has been collected —"
, " so a [retired] line here does not mean that graph is still held in memory."
, ""
, " A line typed at this loop shows nothing more. One run from a `load`ed file ends in"
, " [loaded <file>]; one made by the fetcher (`run serve --follow`) ends in"
, " [fetched <registry> label=<label> id=<document id> sha256=<digest prefix>], which is"
, " how to tell what you typed from what a document said."
, ""
, " The log is capped; if older declarations have fallen off the end, a line after the"
, " listing says how many."
, ""
, " With --select/--exclude, only declarations that named at least one node in the resolved"
, " selection (see `help select`) are shown."
]
queryHelp :: [Text]
queryHelp =
[ "serve: query [--select PATTERN]... [--exclude PATTERN]..."
, ""
, " Lists every node (like a plain `status`), annotating each one [selected] or [excluded]"
, " against the resolved selection (see `help select`), without acting on anything — no"
, " convergence pass runs. Useful for checking what a `converge --select/--exclude` would"
, " touch before actually running it."
]
superviseHelp :: [Text]
superviseHelp =
[ "serve: supervise on|off"
, ""
, " Whether nodes are *tended* between convergence passes, rather than only applied by"
, " them. On by default."
, ""
, " While supervising, every node this world knows has a small state machine of its own,"
, " running in whatever direction the node is wanted. A node wanted up rechecks its own"
, " effect on an adaptive delay — doubling to a minute while the effect is there, halving"
, " to half a second when it is not — and runs `up` again if the effect has gone. A node"
, " wanted down retries its `down` on the same schedule until it succeeds, then stops."
, ""
, " The machines run only while this loop is idle. They start when there is nothing"
, " waiting in the input and stand down before any command is handled (stopping waits for"
, " any `up`/`down` in flight rather than cutting it), so a piped script — every line of"
, " which is queued before the first pass ends — is never supervised at all, and behaves"
, " exactly as it did before any of this existed."
, ""
, " Turning supervision off stops the machines and leaves every effect exactly as it is:"
, " `off` is not a teardown, it is `serve` behaving as it did before nodes had machines."
, ""
, " What a node does about its effect going away is the node author's choice, stated as a"
, " supervision policy on the node itself (see Salmon.Op.Supervision):"
, ""
, " OnFailure put it back when the check says the effect is gone. The default."
, " Never report it and leave it; an operator decides."
, " Always also rerun a node whose check says it ran to completion and stopped."
, ""
, " A check that cannot tell (`Unknown`) never triggers a restart: nothing here re-runs"
, " `up` on a node that looked and could not say. A node with no check of its own — most"
, " of them — answers `Immaterial` instead (\"asking would cost what applying costs\"), is"
, " brought up once, and is then parked rather than polled: still reachable by `force`,"
, " `recheck` and by a dependency that takes its dependants with it, but no longer woken"
, " on a timer to be told the same thing. A node whose effect can go away behind salmon's"
, " back wants a real check; that is what makes it noticeable at all."
, ""
, " A node may instead opt into `supReapply`: rather than being parked, it re-runs `up`"
, " on the same adaptive delay, in place of asking. Only sound for an `up` that is both"
, " cheap and genuinely idempotent — a directory is the case it exists for, a build or a"
, " clone is not — and ignored for a node that owns a running process, whose `up` is not"
, " meant to be re-run at all."
, ""
, " A node may also declare a watchdog: how long it may go without doing anything"
, " observable before that silence should be reported. Nodes that declare none are never"
, " called wedged, which is the right default — silence is evidence only once somebody has"
, " said what silence would mean."
]
autoConvergeHelp :: [Text]
autoConvergeHelp =
[ "serve: autoconverge on|off"
, ""
, " Whether a declaring command (`up`/`only`/`down`/`clear`, and the `-directive` forms)"
, " triggers a convergence pass immediately after recording its epoch. On by default, which"
, " is what makes `up <seed>` on its own bring the seed's nodes up: the declaration and the"
, " pass that acts on it happen as one step."
, ""
, " `autoconverge off` splits that in two. A declaration still updates the active set and"
, " `worldLedger`/`worldNodes` right away — `status`/`query`/`history` see it immediately —"
, " but nothing is applied until an explicit `converge` (optionally restricted with"
, " --select/--exclude). This is the way to record several declarations (e.g. `up a`, then"
, " `down b`, then `up c`) and inspect the combined result with `query`/`status` before"
, " anything actually runs, or to review a directive-driven declaration for a mistake before"
, " committing to it."
, ""
, " Supervision (`help supervise`) is unaffected either way: a node already up and already"
, " supervised keeps being tended regardless of this setting, which only governs whether a"
, " *new* declaration's own pass fires on its own. `converge` (with no autoconverge caveat)"
, " always still runs a pass, whichever way this is set."
]
instructHelp :: [Text]
instructHelp =
[ "serve: force | recheck | pause | resume [--select PATTERN]... [--exclude PATTERN]..."
, ""
, " Tell the matching nodes' own machines something a check cannot: (see `help supervise`"
, " for what those machines are). Omitting --select entirely means every node, same as"
, " `status`/`query`."
, ""
, " force run `up` even though the check says not to — the operator knows something"
, " it does not. For a node that owns a process, this is how to restart one"
, " that is healthy, which is otherwise not sayable at all."
, " recheck look now instead of waiting out the current delay."
, " pause stop tending, without touching the effect. For a node that owns a process,"
, " this leaves it running, unwatched — the operational verb for \"stop caring"
, " about this without stopping it\"."
, " resume start tending again."
, ""
, " These act on a node's machine, and a node only has one while supervision is tending it"
, " (see `help supervise`) — a piped script, or `supervise off`, means there is nothing to"
, " instruct. A one-shot machine (most nodes) does not survive the command that named it"
, " either: `serve` stands every one-shot machine down before handling any command, `status`"
, " included, so there is no live mailbox to post into at the moment this is typed. So the"
, " instruction is queued instead and delivered the moment tending next starts — the next"
, " time this loop goes idle, immediately after the command that queued it. A node a"
, " selection matched that never gets a machine (excluded, retired, or simply never tended)"
, " is not an error: the count this command reports is how many nodes matched, not how many"
, " machines heard it."
]
fetchHelp :: [Text]
fetchHelp =
[ "serve: fetch"
, ""
, " Only meaningful under --follow. The fetcher polls its registry on a schedule: at a base"
, " interval while rounds succeed, backing off (times --follow-factor, up to --follow-cap)"
, " while they fail; and a changed document is not applied at once but held until the"
, " registry has been quiet for --follow-debounce (or --follow-max-wait has elapsed since the"
, " first pending change), so a publisher writing several times in a row is one pass."
, ""
, " `fetch` cuts both short: a round runs now, the backoff is forgotten, and whatever is"
, " pending afterwards is applied without waiting out the quiet window — for an operator"
, " who just published and does not want to wait. Without --follow it only says so."
]
selectHelp :: [Text]
selectHelp =
[ "serve: --select PATTERN / --exclude PATTERN"
, ""
, " Shared by `converge`/`status`/`history`/`query`. A PATTERN is a /-separated glob over a"
, " node's tree position (the same path `run tree`/`run dag` print): a plain segment must"
, " match literally, `*` matches exactly one segment, `**` matches any number of segments"
, " (including zero), so `**` alone matches everything."
, ""
, " Both flags may repeat; each is unioned with itself first. The selected set is every node"
, " matching some --select pattern (or, if --select is omitted entirely, every node), minus"
, " every node matching some --exclude pattern."
, ""
, " The same node (a shared predecessor, e.g. a directory two files sit in) can occur at"
, " several paths; matching any one of them is enough to select or exclude it."
, ""
, " `status`/`query` are where these paths actually come from — each node's listing there"
, " shows every path it currently has, pasteable straight back in as a PATTERN. A path built"
, " from op kinds alone (`directory`, `file-contents`, ...) can be the same for two different"
, " nodes when a recipe reuses the same shorthand at each position; when that happens, prefix"
, " the node's own Ref (also printed on its `status` line) with `#` instead — `#fragment`"
, " matches any node whose Ref starts with that text, which is always unique."
]
-------------------------------------------------------------------------------
{- | Read declarations from a handle until EOF (or @quit@), converging after
each one, and hand back the 'World' as it stands when the loop ends. Never
tears anything down on its way out: exiting the loop leaves the machine as
the last convergence left it. The handle is the loop's standard input
('stdinProducer'); see 'serveProducers' for feeding it from more than one
place.
-}
serve ::
forall seed directive.
(ToJSON directive, FromJSON directive) =>
-- | loop-level events
Reporter Report ->
-- | per-node events, same reporter @run up@ uses
Reporter (UpDown.Report Extension) ->
-- | parses a seed out of one declaration's arguments
([String] -> Either Text seed) ->
Configure IO seed directive ->
Track' directive ->
Handle ->
IO (World seed directive)
serve = serveWith [] Nothing True
{- | 'serve', with "Salmon.Op.Rewrite" phases registered. They run after every
fold, so a convergence walks the /computed/ graph — the one where a
collection node has replaced the nodes it batches — while the ledger and
'worldNodes' keep speaking in terms of what was declared. See
'Salmon.Op.Rewrite' for why that split is the only place cross-declaration
knowledge can live.
The 'Maybe' 'ConcurrencyLimit' bounds each convergence pass's two concurrent
walks (see "Salmon.Actions.Concurrent"): 'Nothing' is unbounded, matching
'serve's behaviour before the limit existed. One limit covers both the
teardown and the bring-up half of every pass, not one each, since the two
never run at the same time (teardown is awaited before bring-up starts) and
so never contend with each other for it.
The 'Bool' is the starting value of @autoconverge@ (see 'AutoConverge'):
'True' matches every version of 'serve' before the setting existed (each
declaration converges immediately), 'False' starts the loop the way an
in-session @autoconverge off@ would, for a caller (e.g. a CLI flag) that
wants declarations held back from the very first line rather than needing
the operator to type it first.
-}
serveWith ::
forall seed directive.
(ToJSON directive, FromJSON directive) =>
[Rewrite Extension] ->
Maybe ConcurrencyLimit ->
Bool ->
Reporter Report ->
Reporter (UpDown.Report Extension) ->
([String] -> Either Text seed) ->
Configure IO seed directive ->
Track' directive ->
Handle ->
IO (World seed directive)
serveWith rewrites limit autoConverge0 r nodeReporter parseSeed configure program h =
serveProducers rewrites limit autoConverge0 r nodeReporter parseSeed configure program [stdinProducer h]
-------------------------------------------------------------------------------
-- input producers
{- | Who typed a line. Standard input is singled out because its end of input
is the one that ends the loop (see 'serveProducers'); every other source is
named, so that a report or a history entry can one day say where a
declaration came from.
-}
data Origin
= -- | the process's own standard input
Stdin
| -- | any other source: a socket connection, a test
Origin !Text
| -- | a line run from a @load@ed file (never pushed by a producer: the
-- loop itself tags the file's lines as it runs them)
Loaded !FilePath
| -- | a declaration "Salmon.Actions.Follow" made from a fetched
-- document; see 'Provenance' for what @history@ says about it
Fetched !Provenance
deriving (Show, Eq, Ord)
{- | Where a fetched declaration came from, in enough detail that an operator
reading @history@ can tell "I typed this" from "the document said so", and
/which/ document: the registry, the label addressed in it, the document's
own @id@ and the digest of its bytes.
-}
data Provenance = Provenance
{ provRegistry :: !Text
, provLabel :: !Text
, provDocument :: !Text
, provDigest :: !Text
}
deriving (Show, Eq, Ord)
{- | An 'Origin' as a report names it in a sentence ("stdin hung up"), as
opposed to 'renderOrigin', the bracketed annotation @history@ appends.
-}
originName :: Origin -> Text
originName Stdin = "stdin"
originName (Origin t) = t
originName (Loaded path) = "loaded " <> Text.pack path
originName (Fetched prov) = "fetched " <> prov.provRegistry <> " label=" <> prov.provLabel
{- | A report, and the 'Origin' of the command it was emitted for: 'Nothing'
for one emitted between commands (the tending loop's), or before the first
and after the last. See 'serveAttributed'.
-}
data Attributed a = Attributed
{ attributedTo :: !(Maybe Origin)
, attributed :: !a
}
deriving (Show, Functor)
{- | What a 'Producer' pushes into the loop's inbox.
A 'Batch' is the unit a fetched document is injected as: its commands run
back to back with @autoconverge@ held off, so the declarations record without
each one converging on its own, then the setting is put back to whatever it
was — an operator's @autoconverge off@ is not silently re-enabled — and one
@converge@ runs. The batch carries its own commands rather than text lines
so a seed's words survive without a quoting round trip, and it is one inbox
entry rather than several so nothing another producer types can land in the
middle of it. Each command carries its own 'Origin', because one batch can
carry several documents' worth of declarations (several labels changed
inside one quiet window) and @history@ must still say which document each
came from.
-}
data Line
= -- | one line of the input language, as the producer read it
Line !Origin !String
| -- | several commands, handled as one: see above
Batch ![(Origin, ServeCommand)]
| -- | this producer has nothing more to say and its thread is about to end
Eof !Origin
deriving (Show, Eq)
{- | A source of 'Line's. 'serveProducers' runs 'produceInto' on a thread of
its own, hands it the loop's one inbox, and kills the thread when the loop
ends; a producer is expected to push an 'Eof' as its last word and return.
-}
newtype Producer = Producer {produceInto :: TChan Line -> IO ()}
-- | Read a handle line by line until end of file, then 'Eof'.
handleProducer :: Origin -> Handle -> Producer
handleProducer origin h = Producer go
where
go inbox = do
eof <- hIsEOF h
if eof
then atomically (writeTChan inbox (Eof origin))
else do
line <- hGetLine h
atomically (writeTChan inbox (Line origin line))
go inbox
{- | The producer 'serveWith' runs: 'handleProducer' with the 'Stdin' origin,
which is what makes its end of input the loop's. The handle need not be the
process's actual standard input — a test's pipe or script file is the same
thing to the loop.
-}
stdinProducer :: Handle -> Producer
stdinProducer = handleProducer Stdin
{- | 'serveWith', fed by any number of 'Producer's rather than one handle.
Each runs on its own thread so that the loop is never itself blocked in a
read: the supervisor's machines run while it waits, and stopping them has to
be able to interleave with a command arriving.
The loop ends on @quit@, or on 'Eof' from the 'Stdin' origin; an 'Eof' from
any other origin is read past. With no 'Stdin' producer in the list, only
@quit@ ends it.
-}
serveProducers ::
forall seed directive.
(ToJSON directive, FromJSON directive) =>
[Rewrite Extension] ->
Maybe ConcurrencyLimit ->
Bool ->
Reporter Report ->
Reporter (UpDown.Report Extension) ->
([String] -> Either Text seed) ->
Configure IO seed directive ->
Track' directive ->
[Producer] ->
IO (World seed directive)
serveProducers rewrites limit autoConverge0 r nodeReporter parseSeed configure program =
serveFollowing rewrites limit autoConverge0 r nodeReporter parseSeed configure program Nothing
{- | What the loop knows of a fetcher producer ("Salmon.Actions.Follow"):
how to wake it (what @fetch@ pulls) and which 'Mode' it is in (what @status@
prints). It is a pair of hooks rather than a producer's methods because the
loop reads lines and does not know which producer it has; these two are the
only things it needs of the fetcher, and both are read-only from its side.
-}
data Followed = Followed
{ followedFetch :: IO ()
-- ^ a round now; see 'Salmon.Actions.Follow.Scheduler.poke'
, followedMode :: IO Mode
-- ^ never 'Interactive'
, followedApplied :: IO [AppliedDocument]
-- ^ the document last applied per label, for the status sink
-- ("Salmon.Actions.Serve.StatusSink"); read-only, a plain read of the
-- fetcher's own cell
}
{- | What the fetcher last applied for one label, as the status sink
publishes it: the label, the document's @id@, its sha256 and when it was
injected. Defined here rather than in "Salmon.Actions.Follow" because the
loop's 'Followed' names it and the fetcher imports the loop, not the other
way round. -}
data AppliedDocument = AppliedDocument
{ appliedDocLabel :: !Text
, appliedDocId :: !Text
, appliedDocDigest :: !Text
, appliedDocAt :: !UTCTime
}
deriving (Show, Eq)
instance ToJSON AppliedDocument where
toJSON a = object ["label" .= a.appliedDocLabel, "id" .= a.appliedDocId, "sha256" .= a.appliedDocDigest, "applied" .= a.appliedDocAt]
instance FromJSON AppliedDocument where
parseJSON = withObject "applied document" $ \o ->
AppliedDocument <$> o .: "label" <*> o .: "id" <*> o .: "sha256" <*> o .: "applied"
{- | Which guarantees apply to the world right now, for @status@ (see
@specs/pull-mode.md@, "what this does not solve"). 'Interactive' when nothing
is followed: every declaration was typed, loaded or batched by a client.
'Following' when a fetcher is running and the world is what the registry
last said. 'Replay' when the registry could not be reached at startup and
the fetcher applied its cached document instead — the world is the last
thing this host knew, not necessarily what the registry says now — until a
later round in which every followed label answers. -}
data Mode = Interactive | Replay | Following
deriving (Show, Eq, Ord)
renderMode :: Mode -> Text
renderMode Interactive = "interactive"
renderMode Replay = "replay"
renderMode Following = "following"
{- | 'serveProducers', with a 'Followed' for what @fetch@ and @status@ ask
of a fetcher producer. 'Nothing' when nothing is being followed: @fetch@
then only says so, and @status@ reports 'Interactive'. -}
serveFollowing ::
forall seed directive.
(ToJSON directive, FromJSON directive) =>
[Rewrite Extension] ->
Maybe ConcurrencyLimit ->
Bool ->
Reporter Report ->
Reporter (UpDown.Report Extension) ->
([String] -> Either Text seed) ->
Configure IO seed directive ->
Track' directive ->
Maybe Followed ->
[Producer] ->
IO (World seed directive)
serveFollowing rewrites limit autoConverge0 r nodeReporter =
serveAttributed rewrites limit autoConverge0 (contramap attributed r) (contramap attributed nodeReporter)
{- | 'serveProducers', reporting through reporters that are told whose
report each one is.
The loop keeps one private "line being handled" cell, written when a line
is taken off the inbox and cleared when its command is done, and every
report — the loop's own and the per-node ones a pass emits — is stamped
with it on the way out ('Salmon.Reporter.pulls'). Nothing else about the
loop changes: a command is handled whole before the next is read, so the
cell is stable for as long as a command's reports are being emitted, and
the concurrent walks a pass runs all report inside that window. What
arrives outside it — a machine tending a node between commands — is stamped
'Nothing'.
The cell is written /after/ 'stopTending', not before: the machines standing
down are not something the operator who typed the command asked for, so
whatever they say on their way out is nobody's.
-}
serveAttributed ::
forall seed directive.
(ToJSON directive, FromJSON directive) =>
[Rewrite Extension] ->
Maybe ConcurrencyLimit ->
Bool ->
Reporter (Attributed Report) ->
Reporter (Attributed (UpDown.Report Extension)) ->
([String] -> Either Text seed) ->
Configure IO seed directive ->
Track' directive ->
Maybe Followed ->
[Producer] ->
IO (World seed directive)
serveAttributed = serveObserved (const (pure ()))
{- | 'serveAttributed', handing an observer a way to read the 'World' before
the first line is read.
The accessor is a plain read of the loop's own cell — never a copy, never a
lock — so what it returns is whatever the loop has committed so far: a
declaration's nodes the moment it is recorded (a pass has not necessarily
run), and the tending snapshot 'stopTending' last filed on each node. It is
what a server answering reads ("Salmon.Actions.Serve.Http") holds instead of
a seat in the inbox, which is the whole of how a read stays a read: it never
stands the machines down and never waits behind a command, including one
whose @up@ is taking a while.
The observer is called once, synchronously, before any producer starts; a
server that wants to run for the loop's lifetime forks from it. The loop
does not kill anything the observer started — a server's own bracket owns
that — but it does return, so an observer holding the accessor after that
reads the final 'World', the same value this returns. It is the first
argument, ahead of everything 'serveAttributed' takes, so that the two
signatures read as one prefixed by the other.
-}
serveObserved ::
forall seed directive.
(ToJSON directive, FromJSON directive) =>
(IO (World seed directive) -> IO ()) ->
[Rewrite Extension] ->
Maybe ConcurrencyLimit ->
Bool ->
Reporter (Attributed Report) ->
Reporter (Attributed (UpDown.Report Extension)) ->
([String] -> Either Text seed) ->
Configure IO seed directive ->
Track' directive ->
Maybe Followed ->
[Producer] ->
IO (World seed directive)
serveObserved observe rewrites limit autoConverge0 rAttributed nodeReporterAttributed parseSeed configure program onFetch producers = do
handling <- newIORef Nothing
serveLoop observe rewrites limit autoConverge0 handling (stamp handling rAttributed) (stamp handling nodeReporterAttributed) parseSeed configure program onFetch producers
where
stamp :: IORef (Maybe Origin) -> Reporter (Attributed a) -> Reporter a
stamp handling = pulls (\rep -> (`Attributed` rep) <$> readIORef handling)
-- | The loop itself: 'serveAttributed' with the stamping already applied
-- and the cell it reads from in hand.
serveLoop ::
forall seed directive.
(ToJSON directive, FromJSON directive) =>
(IO (World seed directive) -> IO ()) ->
[Rewrite Extension] ->
Maybe ConcurrencyLimit ->
Bool ->
IORef (Maybe Origin) ->
Reporter Report ->
Reporter (UpDown.Report Extension) ->
([String] -> Either Text seed) ->
Configure IO seed directive ->
Track' directive ->
Maybe Followed ->
[Producer] ->
IO (World seed directive)
serveLoop observe rewrites limit autoConverge0 handling r nodeReporter parseSeed configure program onFetch producers = do
world <- newIORef emptyWorld
observe (readIORef world)
tending <- Tending <$> newIORef Nothing <*> newIORef Upkeep.noKept <*> newIORef True <*> newIORef autoConverge0 <*> newIORef Map.empty
inbox <- newTChanIO
readers <- traverse (\p -> forkIO (produceInto p inbox)) producers
runReporter r Started
loop tending world inbox `finally` (stopTending tending world >> traverse_ killThread readers)
readIORef world
where
-- | Deepest chain of nested @load@s allowed, to bound a self-referential
-- (or mutually-referential) load file rather than looping forever.
maxLoadDepth :: Int
maxLoadDepth = 8
loop :: Tending -> IORef (World seed directive) -> TChan Line -> IO ()
loop tending world inbox = do
{- Tend the nodes only while there is genuinely nothing to do.
The reader thread queues input as fast as it arrives, so an empty
inbox means the loop is idle and a non-empty one means the next
command is already waiting. Starting machines only when idle is worth
more than the two lines it costs:
* a piped script behaves exactly as it did before any of this
existed. Every line, end-of-input included, is already queued by
the time the first pass finishes, so nothing is ever tended and
@serve < script@ stays a deterministic sequence of passes;
* there is nothing to race. Starting machines and then stopping
them because a command had been sitting in the queue all along
would mean whether a node got acted on depended on thread
timing.
Which leaves supervision doing exactly what it is for: minding the
nodes while whoever is driving this loop is not saying anything. -}
idle <- atomically (isEmptyTChan inbox)
when idle (startTending tending world)
line <- atomically (readTChan inbox)
case line of
-- another producer hanging up is not a command: nothing is about
-- to act, so the machines are not stood down, and the loop goes
-- back to waiting (they are left running if they were). It is
-- said, though: every line that origin typed has been handled
-- by now, which is what whoever holds its connection waits for.
Eof origin | origin /= Stdin -> do
runReporter r (HungUp origin)
loop tending world inbox
_ -> do
-- a command is about to act on these nodes, so the machines
-- stand down. Waits for anything in flight rather than
-- cutting it.
stopTending tending world
case line of
Eof _ -> runReporter r Stopped
Line origin l -> do
writeIORef handling (Just origin)
keepGoing <- step tending 0 world origin l
writeIORef handling Nothing
when keepGoing (loop tending world inbox)
Batch cmds -> do
keepGoing <- batch tending world cmds
when keepGoing (loop tending world inbox)
{- | Run a 'Batch': every command with @autoconverge@ held off, the
setting put back afterwards (a @finally@, so a command that stops the
loop still leaves it as the operator had it), then one full convergence
pass — the sequence the fetcher would otherwise have to spell as
@autoconverge off@ … @autoconverge on@ … @converge@ on the inbox, except
that only the loop knows what to put the setting back /to/. An empty
batch converges nothing: there is no declaration to act on. -}
batch :: Tending -> IORef (World seed directive) -> [(Origin, ServeCommand)] -> IO Bool
batch tending world cmds = do
was <- readIORef (tendingAutoConverge tending)
writeIORef (tendingAutoConverge tending) False
keepGoing <-
runAll cmds `finally` writeIORef (tendingAutoConverge tending) was
when (keepGoing && not (null cmds)) (converge tending world Nothing)
pure keepGoing
where
runAll [] = pure True
runAll ((origin, cmd) : rest) = do
go <- stepCommand tending 0 world origin cmd
if go then runAll rest else pure False
-------------------------------------------------------------------------
-- supervision
{- | Start tending every node this world knows, in whatever direction it
is wanted. Called by 'loop' when it has nothing to do, and stopped again
the moment it has. -}
startTending :: Tending -> IORef (World seed directive) -> IO ()
startTending tending world = do
already <- readIORef (tendingSup tending)
case already of
-- a read-only command does not stand the machines down, so by
-- the time the loop is idle again they are still running and
-- there is nothing to do. Restarting them would re-check every
-- node for no reason, and would drop the mailboxes an operator
-- may have posted into.
Just _ -> pure ()
Nothing -> do
on <- readIORef (tendingOn tending)
auto <- readIORef (tendingAutoConverge tending)
w <- readIORef world
unless (not on || Map.null w.worldNodes) $ do
-- the same computed dag a pass walks: a rewrite's
-- collection node is what actually gets tended, and
-- 'membersOf' is what keeps the bookkeeping in declared
-- terms.
let computed = Rewrite.rewrite rewrites (phaseOf w Nothing) (worldDag w)
kept <- readIORef (tendingKept tending)
forced <- Map.keysSet <$> readIORef (tendingPending tending)
sup <-
Upkeep.startUpkeep
(tendReporter world computed)
kept
(tendOf auto forced w computed)
(Rewrite.computedDag computed)
-- the supervisor owns them now: it adopted what it could
-- and released the rest.
writeIORef (tendingKept tending) Upkeep.noKept
writeIORef (tendingSup tending) (Just sup)
deliverPending tending sup
{- | (R2). Hand every queued instruction to the machine it was meant for,
now that one exists, and forget it. Delivered in the order they were
posted, which matters for e.g. a @pause@ followed by a @resume@.
A node an instruction named that this supervisor is not tending at all
(excluded by the selection at declare time, retired, or simply never
matched a live node) silently drops it here exactly as 'Upkeep.instruct'
always has — there was nothing to queue it *for* once its target never
showed up, and the operator already saw how many nodes matched when the
command was typed ('Instructed'). -}
deliverPending :: Tending -> Upkeep.Supervisor Extension -> IO ()
deliverPending tending sup = do
pending <- readIORef (tendingPending tending)
unless (Map.null pending) $ do
forM_ (Map.toList pending) $ \(aref, instrs) ->
forM_ instrs (Upkeep.instruct sup aref)
writeIORef (tendingPending tending) Map.empty
-- | (R2). Queue an instruction for every named node, oldest first per
-- node, for 'deliverPending' to hand to the next supervisor.
queueInstruction :: Tending -> Set Ref -> Mailbox.Instruction -> IO ()
queueInstruction tending refs instr =
modifyIORef' (tendingPending tending) $ \pending ->
Set.foldr (\aref -> Map.insertWith (flip (<>)) aref [instr]) pending refs
{- | Stop tending, without tearing anything down.
Two different things happen to the two kinds of machine, and the
difference is the whole of why 'Upkeep.Kept' exists. A machine tending an
effect that persists on its own is wound down, waiting for any @up@ or
@down@ in flight rather than interrupting it. A machine /holding/ an
effect up keeps running: this is called before every command, @status@
included, and a supervisor that took its processes with it would restart
every service every time anybody typed anything.
(R3). Before the supervisor's 'TVar's go out of reach, every machine's
'Status' is read and stored on its node — the only place this is ever
readable from, since a running machine's own 'TVar' is not part of
'World' and a discarded 'Upkeep.Supervisor' offers no way back in. Read
from @sup@ itself rather than from 'Upkeep.stopUpkeep''s result, so a
holding machine's status is captured here too and not only a one-shot
one's — 'Upkeep.supervisorStatuses' covers every machine this supervisor
had, before 'Upkeep.stopUpkeep' partitions them into stopped and kept. -}
stopTending :: Tending -> IORef (World seed directive) -> IO ()
stopTending tending world = do
current <- readIORef (tendingSup tending)
forM_ current $ \sup -> do
snapshotStatuses world (Upkeep.supervisorStatuses sup)
kept <- Upkeep.stopUpkeep sup
writeIORef (tendingKept tending) kept
writeIORef (tendingSup tending) Nothing
-- | Read every machine's live 'Status' and file it on its node. See
-- 'stopTending'.
snapshotStatuses :: IORef (World seed directive) -> Map Ref (TVar MachineStatus.Status) -> IO ()
snapshotStatuses world statuses = do
snapshot <- traverse MachineStatus.readStatus statuses
modifyIORef' world $ \w ->
w{worldNodes = Map.foldrWithKey record w.worldNodes snapshot}
where
record aref st = Map.adjust (\ns -> ns{nodeStatus = Just st}) aref
{- | Tear down the machines still holding effects for nodes this world no
longer wants up, and record those nodes as down.
This runs __before__ a convergence pass rather than as part of it, and
the ordering is the point: a daemon's dependencies — its config file, its
working directory — must not be removed while it is still running, and
the down pass is what removes them. 'Upkeep.releaseKept' cancels and
waits, so by the time the pass starts the processes really are gone.
Recording them down here is exact rather than optimistic: for a node
whose effect only exists while something holds it, "nothing holds it" is
what being down /is/. Which is also why a managed node is invisible to
both passes ('gateFor'): there is nothing for a one-shot @down@ to do
that this has not already done, and nothing a one-shot @up@ could do at
all. -}
settleManaged :: Tending -> IORef (World seed directive) -> IO ()
settleManaged tending world = do
w <- readIORef world
kept <- readIORef (tendingKept tending)
kept' <- Upkeep.releaseKept releaseReporter (wantedUp w) kept
writeIORef (tendingKept tending) kept'
let goners =
[ rf
| (rf, a) <- Map.toList w.worldMagma
, isJust a.extension.managed
, Just st <- [Map.lookup rf w.worldNodes]
, st.nodeDirection == TurnDown
]
unless (null goners) $
modifyIORef' world $ \w0 ->
foldr (\rf acc -> setConvergence TurnDown rf Converged acc) w0 goners
wantedUp :: World seed directive -> Ref -> Bool
wantedUp w rf =
case Map.lookup rf w.worldNodes of
Just st -> st.nodeDirection == TurnUp
Nothing -> False
-- | Just enough of 'tendReporter' for 'Upkeep.releaseKept', which is
-- called outside any particular pass and so has no 'Rewritten' to
-- translate through.
releaseReporter :: Reporter (Upkeep.Report Extension)
releaseReporter = ReporterM $ \rep ->
case rep of
Upkeep.Acted inner -> runReporter nodeReporter inner
_ -> runReporter r (Tended rep)
{- | Which nodes the supervisor tends, and how.
'gateFor' is the convergence version of this and differs in one place: it
demands the node has /not/ converged yet, because a pass is one attempt
at whatever is outstanding. Tending is the opposite — a converged node is
precisely the one worth keeping an eye on — so convergence becomes
'Upkeep.Standing' rather than a filter: a converged node starts already
where it wants to be and is only watched, and a 'Pending'\/'Errored'\/
'Blocked' one is acted on.
That distinction is load-bearing rather than an optimisation. Almost no
node in this repository has a @check@, so almost every node answers
'UpDown.Unknown'; without it, starting a supervisor after a pass would
re-run every @up@ in the graph.
The first 'Bool' is @autoconverge@'s current value, and it narrows
"acted on" for a plain (non-'managed') node: with autoconverge off, a
not-yet-'Converged' one-shot node is left untended (returns 'Nothing')
rather than 'Unsettled', because applying it is exactly the convergence
work an operator just asked to defer to an explicit @converge@ — without
this, a node the idle loop reached before that @converge@ would get 'up'
run on it anyway, since almost every node's @check@ answers
'UpDown.Immaterial'\/'UpDown.Unknown' and 'Unsettled' treats either as
"go ahead". A 'managed' node is exempt: it has no other path to ever
start (the convergence pass ignores it categorically, see
'settleManaged'), so autoconverge being off must not also mean "never".
Already-'Converged' nodes are unaffected either way — self-healing an
effect already brought up is not the convergence work being deferred.
The 'Set' 'Ref' is every node with an instruction still queued in
'tendingPending' — a @force@\/@recheck@\/@pause@\/@resume@ typed while
autoconverge is off. Deferring convergence must not also swallow an
operator naming a node explicitly: that instruction has nowhere to be
delivered at all (no machine exists to post it to, see
'deliverPending') unless a machine starts for it here, autoconverge or
not. This is the same exemption 'managed' gets and for the same reason
— an explicit, targeted ask is not the batched convergence work
@autoconverge off@ defers — it just reaches that ask through a
different field than @managed@ does.
-}
tendOf :: Bool -> Set Ref -> World seed directive -> Rewritten Extension -> Ref -> Maybe Upkeep.Tend
tendOf autoConverge forced w computed aref =
case [st | rf <- Set.toList (Rewrite.membersOf computed aref), Just st <- [Map.lookup rf w.worldNodes]] of
[] -> Nothing
sts ->
-- a collection node standing in for members that disagree
-- goes up: the conservative direction, the same call
-- 'Salmon.Op.Rewrite' asks its phases to make. And it counts
-- as standing only if /every/ member it speaks for does,
-- which is the same all-or-nothing attribution a batch makes
-- everywhere else.
let ups = [st | st <- sts, st.nodeDirection == TurnUp]
mine = if null ups then sts else ups
converged = all (\st -> st.nodeConvergence == Converged) mine
managed = maybe False (isJust . (.extension.managed)) (Map.lookup aref (Dag.dagNodes (Rewrite.computedDag computed)))
instructed = not (Set.null (Set.intersection forced (Rewrite.membersOf computed aref)))
in if not converged && not autoConverge && not managed && not instructed
then Nothing
else
Just
Upkeep.Tend
{ Upkeep.tendDirection = if null ups then TurnDown else TurnUp
, Upkeep.tendStanding =
if converged
then Upkeep.Settled
else Upkeep.Unsettled
}
{- | Where a machine's reports go.
Node-level events ('Upkeep.Acted') are the one-shot drivers' own
vocabulary, so they go where a pass's do: into the convergence
bookkeeping, and on to the caller's node reporter. Two filters, both
about volume rather than meaning:
* a 'UpDown.Skip' is recorded but not printed. The supervisor re-checks
every node when it starts, and saying "nothing to do" once per node
per convergence on top of what the pass already said is noise;
* 'Upkeep.NextLook' and the state transitions are dropped entirely. Every
node emits one on every nap, forever, which is a trace rather than a
report. What survives is what an operator would want woken for: a
wedged node, a paused one, a contradictory policy, a machine that
escaped. -}
tendReporter :: IORef (World seed directive) -> Rewritten Extension -> Reporter (Upkeep.Report Extension)
tendReporter world computed = ReporterM $ \rep ->
case rep of
Upkeep.Acted inner -> do
runReporter (tendWriter world computed) inner
case inner of
UpDown.Skip _ -> pure ()
_ -> runReporter nodeReporter inner
Upkeep.Upkeep{} -> pure ()
Upkeep.Downkeep{} -> pure ()
Upkeep.NextLook{} -> pure ()
Upkeep.Untended{} -> pure ()
_ -> runReporter r (Tended rep)
{- | 'stateWriter', for a driver that tends both directions at once.
The convergence version is told which direction its pass is for; a
supervisor is not, so each node's own currently-wanted direction is what
its outcome is recorded against. There is no @restriction@ either: a
supervisor is never scoped by a @--select@, because the operator restricts
a /pass/, not what is kept running. -}
tendWriter :: IORef (World seed directive) -> Rewritten Extension -> Reporter (UpDown.Report Extension)
tendWriter world computed = ReporterM $ \rep ->
case rep of
UpDown.Eval _ -> pure ()
UpDown.Done act -> mark act Converged
UpDown.Skip act -> mark act Converged
UpDown.Failed act _ -> mark act Errored
UpDown.Blocked act -> mark act Blocked
UpDown.Conflicting{} -> pure ()
UpDown.Instructed{} -> pure ()
UpDown.DroppedInstructions{} -> pure ()
where
mark :: Act Extension -> Convergence -> IO ()
mark act c =
forM_ (Set.toList (Rewrite.membersOf computed act.extension.ref)) $ \rf ->
atomicModifyIORef' world (\w -> (setConvergenceHere rf c w, ()))
step :: Tending -> Int -> IORef (World seed directive) -> Origin -> String -> IO Bool
step tending depth world origin line =
case parseServeCommand line of
Left err -> do
runReporter r (BadCommand err)
pure True
Right cmd -> stepCommand tending depth world origin cmd
stepCommand :: Tending -> Int -> IORef (World seed directive) -> Origin -> ServeCommand -> IO Bool
stepCommand tending depth world origin cmd =
case cmd of
Noop -> pure True
Quit -> pure False
Help mtopic -> do
runReporter r (HelpText mtopic)
pure True
Status sel -> do
w <- readIORef world
mode <- maybe (pure Interactive) followedMode onFetch
runReporter r (StatusReport mode (filterNodes w sel) (worldPaths w))
pure True
History sel -> do
w <- readIORef world
let (selr, excr) = resolveWorldSelectors w sel
let allowed = selr `Set.difference` excr
let matches :: LogEntry -> Bool
matches e = sel == noSelection || not (Set.null (Set.intersection e.logRefs allowed))
runReporter r (HistoryReport (historyLinesMatching matches w))
when (w.worldLogDropped > 0) $
runReporter r (HistoryElided w.worldLogDropped)
pure True
QueryCmd sel -> do
w <- readIORef world
let (selr, excr) = resolveWorldSelectors w sel
runReporter r (QueryReport (Map.toList w.worldNodes) selr excr (worldPaths w))
pure True
Converge sel -> do
restriction <-
if sel == noSelection
then pure Nothing
else do
w <- readIORef world
let (selr, excr) = resolveWorldSelectors w sel
pure (Just (selr `Set.difference` excr))
converge tending world restriction
pure True
Clear -> do
w <- readIORef world
writeIORef world (resettle w{worldLedger = Ledger.retractAll w.worldLedger})
runReporter r (Cleared (Ledger.liveCount w.worldLedger))
convergeIfAuto tending world
pure True
Declare decl args -> do
declare tending world origin decl args
pure True
DeclareDirective decl path -> do
declareDirective tending world origin decl path
pure True
DeclareInline decl name value -> do
declareDecoded tending world origin decl ["<directive>", Text.unpack name] (parseEither parseJSON value)
pure True
Load path -> loadFile tending world (depth + 1) path
Supervise on -> do
writeIORef (tendingOn tending) on
-- turning it off has to take effect now; turning it
-- on happens the moment this loop is next idle,
-- which is immediately after this command.
unless on (stopTending tending world)
runReporter r (Supervised on)
pure True
AutoConverge on -> do
writeIORef (tendingAutoConverge tending) on
runReporter r (AutoConverged on)
pure True
Instruct instr sel -> do
w <- readIORef world
let (selr, excr) = resolveWorldSelectors w sel
allowed = selr `Set.difference` excr
queueInstruction tending allowed instr
runReporter r (Instructed instr (Set.size allowed))
pure True
Fetch -> do
traverse_ followedFetch onFetch
runReporter r (FetchRequested (isJust onFetch))
pure True
-- | Filters 'worldNodes' by a 'Selection', preserving today's exact
-- unfiltered listing (including nodes wanted 'TurnDown') when no
-- @--select@\/@--exclude@ was given at all.
filterNodes :: World seed directive -> Selection -> [(Ref, NodeState)]
filterNodes w sel
| sel == noSelection = Map.toList w.worldNodes
| otherwise =
let (selr, excr) = resolveWorldSelectors w sel
allowed = selr `Set.difference` excr
in [(rf, st) | (rf, st) <- Map.toList w.worldNodes, rf `Set.member` allowed]
loadFile :: Tending -> IORef (World seed directive) -> Int -> FilePath -> IO Bool
loadFile tending world depth path
| depth > maxLoadDepth = do
runReporter r (BadLoad ("refusing to load " <> Text.pack path <> ": nesting too deep (possible cycle)"))
pure True
| otherwise = do
runReporter r (Loading path)
result <- try (readFile path) :: IO (Either IOException String)
case result of
Left ex -> do
runReporter r (BadLoad ("cannot read " <> Text.pack path <> ": " <> Text.pack (show ex)))
pure True
Right contents -> go 0 (lines contents)
where
go n [] = do
runReporter r (LoadDone path n)
pure True
go n (ln : rest) = do
keepGoing <- step tending depth world (Loaded path) ln
if keepGoing then go (n + 1) rest else pure False
-- | A 'Configure' that throws is reported as a bad seed and the loop
-- reads on, same as a seed that fails to parse. It used to take the whole
-- loop down, which for a typed line was a nuisance and for a fetched
-- document (whose author is not at this keyboard) would be a host
-- losing its supervisor to somebody else's typo.
declare :: Tending -> IORef (World seed directive) -> Origin -> Declaration -> [String] -> IO ()
declare tending world origin decl args =
case parseSeed args of
Left err -> runReporter r (BadSeed err)
Right seed -> do
configured <- try (gen configure seed) :: IO (Either SomeException directive)
case configured of
Left ex -> runReporter r (BadSeed (Text.pack (unwords args) <> ": configure threw: " <> Text.pack (show ex)))
Right directive -> declareConfigured tending world origin decl args seed directive
declareConfigured :: Tending -> IORef (World seed directive) -> Origin -> Declaration -> [String] -> seed -> directive -> IO ()
declareConfigured tending world origin decl args seed directive = do
w0 <- readIORef world
let o = run program directive
let gr = evalDeps o
let ep =
Epoch
{ epochId = EpochId w0.worldNextId
, epochDeclaration = decl
, epochDirection = declarationDirection decl
, epochOrigin = origin
, epochTokens = args
, epochSeed = Just seed
, epochDirective = directive
, epochKey = encode directive
, epochGraph = gr
}
commitEpoch tending world w0 decl ep
declareDirective :: Tending -> IORef (World seed directive) -> Origin -> Declaration -> FilePath -> IO ()
declareDirective tending world origin decl path = do
result <- try (LByteString.readFile path) :: IO (Either IOException ByteString)
case result of
Left ex -> runReporter r (BadDirective ("cannot read " <> Text.pack path <> ": " <> Text.pack (show ex)))
Right bytes -> declareDecoded tending world origin decl ["<directive-file>", path] (eitherDecode bytes)
-- | The tail of a directive declaration once its JSON has been read from
-- wherever it was: a decode failure is reported and nothing is declared.
declareDecoded :: Tending -> IORef (World seed directive) -> Origin -> Declaration -> [String] -> Either String directive -> IO ()
declareDecoded tending world origin decl tokens decoded =
case decoded of
Left err -> runReporter r (BadDirective (Text.pack err))
Right directive -> do
w0 <- readIORef world
let o = run program directive
let gr = evalDeps o
let ep =
Epoch
{ epochId = EpochId w0.worldNextId
, epochDeclaration = decl
, epochDirection = declarationDirection decl
, epochOrigin = origin
, epochTokens = tokens
, epochSeed = Nothing
, epochDirective = directive
, epochKey = encode directive
, epochGraph = gr
}
commitEpoch tending world w0 decl ep
-- | Appends and records a freshly-built epoch, then converges (fully:
-- a declaration is never itself scoped by a 'Selection') — unless
-- @autoconverge off@ has asked declarations to just record and wait.
commitEpoch :: Tending -> IORef (World seed directive) -> World seed directive -> Declaration -> Epoch seed directive -> IO ()
commitEpoch tending world w0 decl ep = do
let dag = Dag.foldDag Dag.sameRepresentative ep.epochGraph
-- the fold is where a Ref collision inside one declaration is
-- visible; the down pass no longer folds anything, so this is the
-- only place left that can say so.
forM_ (reverse (Dag.dagConflicts dag)) $ \c ->
runReporter nodeReporter (UpDown.Conflicting c.conflictRef c.conflictKept c.conflictReplaced)
let (recorded, crossed) = recordWith decl ep dag w0
-- ... and a collision with another live declaration is only
-- visible once the ledger says who else wants the node.
forM_ crossed $ \c ->
runReporter nodeReporter (UpDown.Conflicting c.conflictRef c.conflictKept c.conflictReplaced)
let w1 = resettle recorded
writeIORef world w1
runReporter r $
Declared
ep.epochId
ep.epochDirection
(Map.size (Dag.dagNodes dag))
(length w1.worldEpochs)
convergeIfAuto tending world
-- | 'converge's the whole world, unless @autoconverge off@ is in
-- effect, in which case a declaring command's own report is the only
-- thing the operator sees until an explicit @converge@.
convergeIfAuto :: Tending -> IORef (World seed directive) -> IO ()
convergeIfAuto tending world = do
auto <- readIORef (tendingAutoConverge tending)
when auto (converge tending world Nothing)
-- | Runs one down-then-up convergence pass. @restriction@, when
-- present, additionally 'Skippable'-gates any node whose 'Ref' isn't in
-- it — used only by an explicit @converge --select\/--exclude@; the
-- auto-converge that follows every declaration always passes 'Nothing'.
converge :: Tending -> IORef (World seed directive) -> Maybe (Set Ref) -> IO ()
converge tending world restriction = do
-- 'loop' has already stood the one-shot machines down. What it did
-- not do is let go of the effects something is still /holding/,
-- because at that point this world had not yet been told what the
-- command changed. Now it has, so: anything no longer wanted up goes
-- first, before the down pass starts removing what it stood on.
settleManaged tending world
w <- readIORef world
-- the rewrites run per pass rather than per declaration, because
-- what they partition on ('Ledger.desired') is a property of the
-- whole ledger at this moment, not of any one declaration.
let computed = Rewrite.rewrite rewrites (phaseOf w restriction) (worldDag w)
let dag = Rewrite.computedDag computed
let (nup, ndown) = pendingCounts w
runReporter r (ConvergeStart ndown nup)
-- teardown first: a node being replaced by an incompatible one
-- (different content, hence a different 'Ref') has to go before its
-- successor is brought up.
okDown <-
if ndown == 0
then pure True
else
Concurrent.downDagConcurrent
(gateFor world computed TurnDown restriction)
(recorder world computed TurnDown restriction)
Concurrent.noMailboxes
limit
dag
okUp <-
if nup == 0
then pure True
else
Concurrent.upDagConcurrent
(gateFor world computed TurnUp restriction)
(recorder world computed TurnUp restriction)
Concurrent.noMailboxes
limit
dag
-- this pass is what turns nodes converged-'TurnDown', so it is also
-- where the graphs that described them stop being needed.
modifyIORef' world resettle
w' <- readIORef world
let (rup, rdown) = pendingCounts w'
runReporter r (ConvergeStop (okDown && okUp) (rup + rdown))
{- | Only touch what this pass is for: a node wanted the other way (it
belongs to some other seed), already converged, or excluded by this
pass's own 'restriction' (an explicit @converge --select\/--exclude@) is
left alone. -}
gateFor :: IORef (World seed directive) -> Rewritten Extension -> Direction -> Maybe (Set Ref) -> UpDown.Gate Extension
gateFor world computed dir restriction = \act -> do
w <- readIORef world
-- a node a rewrite introduced has no 'NodeState' of its own; it is
-- worth touching iff any of the declared nodes it stands in for is.
-- For every other node 'membersOf' is the singleton of itself, so
-- this is the same predicate it always was.
pure $
-- a node whose effect only exists while something holds it is
-- not this pass's business in either direction: bringing it up
-- needs a driver that can hold it (so its @up@ throws, on
-- purpose — see "Salmon.Builtin.Nodes.Daemon"), and taking it
-- down is 'settleManaged', which has already run.
if isJust act.extension.managed
then Skippable
else
if any (wants w) (Set.toList (Rewrite.membersOf computed act.extension.ref))
then Required
else Skippable
where
wants :: World seed directive -> Ref -> Bool
wants w rf =
case Map.lookup rf w.worldNodes of
Nothing -> False
Just st ->
st.nodeDirection == dir
&& st.nodeConvergence /= Converged
&& maybe True (Set.member rf) restriction
recorder :: IORef (World seed directive) -> Rewritten Extension -> Direction -> Maybe (Set Ref) -> Reporter (UpDown.Report Extension)
recorder world computed dir restriction = reportBoth (stateWriter world computed dir restriction) nodeReporter
{- 'upTree'/'downTree' report an 'Eval' before running a node and, once
it returns, exactly one of 'Done' (succeeded) or 'Failed' (threw) — so
recording convergence off 'Done'/'Failed' rather than 'Eval' is exact. A
'Skip' is either this pass's own gate (already converged, not ours — both
fine to record as converged, the direction check below drops the latter
— or restricted out by an explicit @converge --select\/--exclude@, which
must leave the node's actual convergence untouched so a later
unrestricted @converge@ still retries it) or, on the way up, the node's
own 'check' saying its effect is already in place, which is convergence
too. -}
stateWriter :: IORef (World seed directive) -> Rewritten Extension -> Direction -> Maybe (Set Ref) -> Reporter (UpDown.Report Extension)
stateWriter world computed dir restriction = ReporterM $ \rep ->
case rep of
UpDown.Eval _ -> pure ()
UpDown.Done act -> mark act Converged
UpDown.Skip act
| maybe False (Set.notMember act.extension.ref) restriction -> pure ()
-- 'gateFor' skips every managed node, and that skip says
-- nothing about whether the node is up: only the machine
-- holding it can say that, and it does so through
-- 'tendWriter'.
| isJust act.extension.managed -> pure ()
| otherwise -> mark act Converged
UpDown.Failed act _ -> mark act Errored
UpDown.Blocked act -> mark act Blocked
-- not a node outcome: it says two declarations describe one
-- node differently, which the operator wants to see but which
-- leaves no node any more or less converged than it was.
UpDown.Conflicting{} -> pure ()
-- likewise not node outcomes: an instruction being applied, or
-- an older one being evicted, says what was asked for rather
-- than what happened.
UpDown.Instructed{} -> pure ()
UpDown.DroppedInstructions{} -> pure ()
where
-- what happened to a collection node happened to every declared node
-- it stands in for — that is the whole of what makes a batch's
-- outcome legible in per-package terms, and it is why a batch
-- reports failure for all of its members.
-- 'atomicModifyIORef'', not 'modifyIORef'': the concurrent driver
-- runs several nodes at once and they all report into this same
-- world, so a read-modify-write that is not atomic silently loses
-- convergence records.
mark :: Act Extension -> Convergence -> IO ()
mark act c =
forM_ (Set.toList (Rewrite.membersOf computed act.extension.ref)) $ \rf ->
atomicModifyIORef' world (\w -> (setConvergence dir rf c w, ()))
-------------------------------------------------------------------------------
{- | Folds one declaration in: its nodes into 'worldMagma', its nodes and
edges into 'worldLedger', its graph into 'worldEpochs', and a line into
'worldLog'. 'worldNextId' only ever grows, so an id in the log stays
meaningful after 'prune' has collected the epoch it names.
Every declaration is 'Ledger.declare'd before the retraction is applied,
@down@ included. That is not a detour: a @down@ re-evaluates its seed, and
folding that evaluation in first is what makes the teardown use the /current/
description of those nodes rather than whatever was declared last time. It is
also why the ledger entry is replaced rather than accumulated — the same key
declared twice is one declaration, so one @down@ retracts it.
-}
record :: Declaration -> Epoch seed directive -> Dag Extension -> World seed directive -> World seed directive
record decl ep dag w = fst (recordWith decl ep dag w)
{- | 'record', also handing back the collisions this declaration has with
/other/ live declarations (one per 'Ref', kept-and-replaced), which the loop
reports 'UpDown.Conflicting' beside the ones the fold found inside the
declaration itself. See 'Collision' for the rule.
-}
recordWith :: Declaration -> Epoch seed directive -> Dag Extension -> World seed directive -> (World seed directive, [Dag.Conflict Extension])
recordWith decl ep dag w =
( w
{ worldNextId = w.worldNextId + 1
, worldEpochs = ep : w.worldEpochs
, worldLog = kept
, worldLogDropped = w.worldLogDropped + length dropped
, -- left-biased: this declaration's representatives win, which is
-- 'Salmon.Op.Dag''s last-writer-wins across declarations.
worldMagma = Map.union (Dag.dagNodes dag) w.worldMagma
, worldLedger = ledger'
, worldConflicts = Map.union collisions (Map.withoutKeys w.worldConflicts described)
, -- (I6): a 'Ref' this declaration redescribes goes 'Stale' rather
-- than staying silently 'Converged' under a representative it was
-- never actually applied against.
worldNodes = foldr demoteIfChanged w.worldNodes (Set.toList changed)
}
, crossed
)
where
contrib = Ledger.contribution dag
described = Map.keysSet (Dag.dagNodes dag)
retraction = case decl of
Add -> id
Replace -> Ledger.retractOthers ep.epochKey
Remove -> Ledger.retract ep.epochKey
ledger' = retraction (Ledger.declare ep.epochKey contrib w.worldLedger)
-- the other live declarations still wanting a node — read off the
-- ledger /after/ the retraction, so an @only@ does not collide with the
-- very seeds it is retiring
othersHolding :: Ref -> Set ByteString
othersHolding rf =
Map.keysSet (Map.filterWithKey (\k c -> k /= ep.epochKey && c.contribLive && Set.member rf c.contribRefs) ledger')
-- the fold's own collisions, oldest first so the newest wins the map
inside :: Map Ref (Dag.Conflict Extension)
inside = Map.fromList [(c.conflictRef, c) | c <- reverse (Dag.dagConflicts dag)]
-- the collisions this declaration is the last writer of: a node it
-- describes differently from the magma while another live declaration
-- still wants it. Reported, as the fold's own are.
crossed :: [Dag.Conflict Extension]
crossed =
[ Dag.Conflict rf newAct oldAct
| (rf, newAct) <- Map.toList (Dag.dagNodes dag)
, Set.member rf changed
, Just oldAct <- [Map.lookup rf w.worldMagma]
, not (Set.null (othersHolding rf))
]
crossedByRef :: Map Ref (Dag.Conflict Extension)
crossedByRef = Map.fromList [(c.conflictRef, c) | c <- crossed]
-- one entry per 'Ref' this declaration describes, or none
collisions :: Map Ref Collision
collisions = Map.mapMaybe id (Map.mapWithKey collisionOf (Dag.dagNodes dag))
collisionOf :: Ref -> Act Extension -> Maybe Collision
collisionOf rf _
| Just c <- Map.lookup rf crossedByRef =
Just (Collision c (othersHolding rf))
| Just c <- Map.lookup rf inside =
Just (Collision c (Set.singleton ep.epochKey))
| not (Set.member rf changed)
, Just standing <- Map.lookup rf w.worldConflicts
, any (`Ledger.isLive` ledger') (Set.toList standing.collisionHolders) =
Just standing
| otherwise = Nothing
where
others = othersHolding rf
{- | Every 'Ref' this declaration describes differently than whatever is
already in the magma — the same 'Dag.sameRepresentative' comparison
'Dag.foldDag' itself uses to decide a re-declaration is a genuine
conflict rather than the overwhelmingly common "one node, reached
again" case. A brand-new 'Ref' (absent from 'worldMagma') is not
"changed": it has nothing to differ from, and 'retune' already gives it
a fresh 'Pending' on its own.
-}
changed :: Set Ref
changed =
Set.fromList
[ rf
| (rf, newAct) <- Map.toList (Dag.dagNodes dag)
, Just oldAct <- [Map.lookup rf w.worldMagma]
, not (Dag.sameRepresentative oldAct newAct)
]
-- only a node currently believed 'Converged' has anything to lose by
-- this: one already 'Pending'\/'Stale'\/'Errored'\/'Blocked' is getting
-- a fresh look regardless, and relabelling it would only blur why.
demoteIfChanged :: Ref -> Map Ref NodeState -> Map Ref NodeState
demoteIfChanged rf =
Map.adjust (\st -> if st.nodeConvergence == Converged then st{nodeConvergence = Stale} else st) rf
entry =
LogEntry
{ logEpoch = ep.epochId
, logDeclaration = ep.epochDeclaration
, logOrigin = ep.epochOrigin
, logTokens = ep.epochTokens
, logRefs = Ledger.contribRefs contrib
}
(kept, dropped) = splitAt worldLogLimit (entry : w.worldLog)
{- | Re-derives the world after anything that could have changed it: first
'retune' (every node's wanted 'Direction', from the active seeds), then
'prune' (drop what is finished).
The order is load-bearing and is the one way to get this wrong. 'prune' asks
which nodes are still on their way down, and immediately after a @down@
declaration is 'record'ed those nodes still look 'TurnUp' and 'Converged' —
so pruning first would collect the very contribution the teardown is about to
be run from. 'retune' is what flips them to 'TurnDown'\/'Pending', after
which 'prune' keeps their contribution. (@Test.ServeSpec@'s "a retired
declaration survives a failed down" case pins this.)
-}
resettle :: World seed directive -> World seed directive
resettle = prune . retune
{- | Re-derives every node's wanted 'Direction' from the active seeds. A node
whose direction is unchanged keeps its 'Convergence'; one that just flipped
goes back to 'Pending', because whatever was done to it was done the other way.
-}
retune :: World seed directive -> World seed directive
retune w =
w{worldNodes = Map.mapWithKey adjust known}
where
desired :: Set Ref
desired = Ledger.desired w.worldLedger
-- every node any retained contribution still mentions, with its metadata
-- taken from the magma — i.e. from the last declaration to describe it.
known :: Map Ref (ShortHand, Text)
known =
Map.fromList
[ (r, (act.shorthand, act.extension.help))
| r <- Set.toList (Ledger.knownRefs w.worldLedger)
, Just act <- [Map.lookup r w.worldMagma]
]
adjust r (sh, hlp) =
let dir = if Set.member r desired then TurnUp else TurnDown
in case Map.lookup r w.worldNodes of
Just st
| st.nodeDirection == dir ->
st{nodeShorthand = sh, nodeHelp = hlp}
_ -> NodeState sh hlp dir Pending Nothing
{- | Drops what is finished, which is what keeps a long-lived @serve@
bounded. Three rules that have to agree with each other:
* a node converged 'TurnDown' is done — it is off the machine and nothing
will be done to it again — so it leaves 'worldNodes' and 'worldMagma';
* a /retired/ contribution is kept only while one of its nodes is still to
be turned down, since its edges are the only remaining statement of what
order to do that in. A live one is never dropped: it is what holds its
nodes up;
* an epoch is kept only while its declaration is live, and then only the
newest for that key. Nothing else needs a graph any more — this is where
the storage saving is, and it is the rule that used to also have to keep
a retired seed's graph for the teardown to walk.
The two 'Map.restrictKeys' are implied by the rules rather than adding to
them: they make "every node in 'worldNodes' is described by some retained
contribution, and every one has a representative" hold structurally instead
of by argument.
Note this is why re-declaring an unchanged seed is cheap forever: the
superseded epoch is no longer newest for its key while its nodes stay
'TurnUp' under the new one, so its graph goes.
-}
prune :: forall seed directive. World seed directive -> World seed directive
prune w =
w
{ worldEpochs = keptEpochs
, worldLedger = ledger
, worldMagma = Map.restrictKeys w.worldMagma (Map.keysSet retained)
, worldConflicts = Map.filter standing (Map.restrictKeys w.worldConflicts (Map.keysSet retained))
, worldNodes = retained
}
where
nodes = Map.filter (not . finished) w.worldNodes
-- a collision stands while somebody on its losing side is still live
standing :: Collision -> Bool
standing c = any (`Ledger.isLive` ledger) (Set.toList c.collisionHolders)
retained = Map.restrictKeys nodes (Ledger.knownRefs ledger)
-- these locals are annotated because a record-dot binding without a
-- signature generalizes over 'HasField' and would need FlexibleContexts.
finished :: NodeState -> Bool
finished st = st.nodeDirection == TurnDown && st.nodeConvergence == Converged
ledger = Ledger.collect stillToTurnDown w.worldLedger
stillToTurnDown :: Ref -> Bool
stillToTurnDown r =
case Map.lookup r nodes of
Just st -> st.nodeDirection == TurnDown
Nothing -> False
-- newest-first, so the first epoch seen for a key is the current one.
keptEpochs :: [Epoch seed directive]
keptEpochs = go Set.empty w.worldEpochs
where
go :: Set ByteString -> [Epoch seed directive] -> [Epoch seed directive]
go _ [] = []
go seen (ep : eps)
| Set.member ep.epochKey seen = go seen eps
| Ledger.isLive ep.epochKey ledger = ep : go (Set.insert ep.epochKey seen) eps
| otherwise = go (Set.insert ep.epochKey seen) eps
setConvergence :: Direction -> Ref -> Convergence -> World seed directive -> World seed directive
setConvergence dir r c w =
w{worldNodes = Map.adjust upd r w.worldNodes}
where
upd st
| st.nodeDirection == dir = st{nodeConvergence = c}
| otherwise = st
{- | 'setConvergence' against whichever direction the node is currently
wanted in, rather than against a stated one.
The convergence passes know their own direction and use it as a filter — a
node wanted the other way belongs to another pass and must not be recorded.
A supervisor tends both directions at once and has no such filter to apply,
so the node's own state is the answer.
-}
setConvergenceHere :: Ref -> Convergence -> World seed directive -> World seed directive
setConvergenceHere r c w =
case Map.lookup r w.worldNodes of
Nothing -> w
Just st -> setConvergence st.nodeDirection r c w
-- | (nodes wanted up, nodes wanted down) that have not converged yet.
pendingCounts :: World seed directive -> (Int, Int)
pendingCounts w =
(count TurnUp, count TurnDown)
where
count dir = length [() | st <- Map.elems w.worldNodes, st.nodeDirection == dir, st.nodeConvergence /= Converged]
{- | What the "Salmon.Op.Rewrite" phases are told about the pass about to
run: which nodes some live declaration still wants (so a rewrite can tell an
install from a removal), and which ones an explicit @converge
--select@\/@--exclude@ has put out of scope (so a rewrite does not quietly
batch up work the operator asked to skip).
-}
phaseOf :: World seed directive -> Maybe (Set Ref) -> Phase
phaseOf w restriction =
Phase
{ phaseDesired = Ledger.desired w.worldLedger
, phaseIgnored = maybe Set.empty (Map.keysSet w.worldNodes `Set.difference`) restriction
}
{- | What both convergence passes walk: the magma, wired back up with the
precedence the ledger holds. No graph is involved, which is the point — a
retired declaration's graph is long gone, and its two flat sets are enough.
One structure for both directions, rather than a union of graphs per pass.
Nodes this pass is not for are in it too, exactly as they used to be in the
epoch graphs the old @upOps@\/@downOps@ handed over; the pass's
'UpDown.Gate' is what leaves them alone, and a 'UpDown.Skip'ped node releases
its neighbours just like an applied one.
-}
worldDag :: World seed directive -> Dag Extension
worldDag w = Dag.fromMagma w.worldMagma (Ledger.precedenceOf w.worldLedger)
{- | Read off 'worldLog', not 'worldEpochs' — a declaration is still worth
printing long after 'prune' has collected the graph it made. The
@[active]@\/@[retired]@ flag stays exact regardless: a collected epoch is
never in 'worldActive'. The only caller is 'HistoryReport', with a predicate
of @const True@ for a plain @history@; there is no unfiltered version left
to call directly, since there was never a caller for one.
-}
historyLinesMatching ::
(LogEntry -> Bool) ->
World seed directive ->
[(EpochId, Declaration, Bool, Origin, [String])]
historyLinesMatching p w =
[ (e.logEpoch, e.logDeclaration, Set.member e.logEpoch activeIds, e.logOrigin, e.logTokens)
| e <- reverse w.worldLog
, p e
]
where
activeIds = activeEpochIds w
{- | Resolves a 'Selection' against every currently-/active/ epoch's graph,
unioning the per-epoch matches — there is no single unified cograph for the
whole 'World', only the unified 'worldNodes' map. An empty 'selSelect' still
resolves to "everything" per epoch, so the union over active epochs is
exactly every active node, mirroring 'retune''s own @desired@ computation.
A pattern beginning with @#@ is, exactly as 'Query.resolveRewrittenSelectors'
already does for @run up@\/@run down@, matched by 'Ref' instead of by path: a
fragment of the text 'status'\/'query' now print on every node's line (either
the short, disambiguating tag or the full 'Ref'). This is what makes a node
addressable at all when two of them share every path — a recipe that reuses
the same shorthand (\"directory\", \"file-contents\", ...) at each position
gives 'Query.pathedRefs' no way to tell them apart by path, and printing the
paths in 'worldPaths' cannot invent a distinction that was never there.
-}
resolveWorldSelectors :: World seed directive -> Selection -> (Set Ref, Set Ref)
resolveWorldSelectors w sel =
(selectedBase `Set.difference` excluded, excluded)
where
(selRefPats, selPathPats) = List.partition isRefFragment sel.selSelect
(excRefPats, excPathPats) = List.partition isRefFragment sel.selExclude
allRefs = Set.unions [Set.fromList (map snd (Query.pathedRefs ep.epochGraph)) | ep <- w.worldEpochs]
pathMatches :: [Text] -> Set Ref
pathMatches [] = Set.empty
pathMatches pats = Set.unions [fst (Query.resolveSelectors ep.epochGraph pats []) | ep <- w.worldEpochs]
refMatches :: [Text] -> Set Ref
refMatches pats = Set.fromList [rf | rf <- Set.toList allRefs, pat <- pats, matchesRefFragment pat rf]
matchesRefFragment :: Text -> Ref -> Bool
matchesRefFragment pat rf =
let fragment = Text.drop 1 pat
in fragment `Text.isPrefixOf` Query.shortRef rf || fragment `Text.isPrefixOf` unRef rf
isRefFragment :: Text -> Bool
isRefFragment = Text.isPrefixOf "#"
matchesOf :: [Text] -> [Text] -> Set Ref
matchesOf pathPats refPats = pathMatches pathPats `Set.union` refMatches refPats
selectedBase = if null sel.selSelect then allRefs else matchesOf selPathPats selRefPats
excluded = matchesOf excPathPats excRefPats
{- | The epochs 'prune' retained are exactly the live declarations' newest
ones, so this needs no separate active-seed index — the ledger's liveness is
the only source of truth for what is declared up.
-}
activeEpochIds :: World seed directive -> Set EpochId
activeEpochIds w = Set.fromList (fmap epochId w.worldEpochs)
{- | Every path (rendered @\/@-separated, root-to-node, exactly the shape
@--select@\/@--exclude@ patterns match against) at which a live declaration's
graph reaches each 'Ref' — the thing @status@\/@query@ never showed despite
being the only practical way to /build/ a selector pattern in the first
place: without this, a node was nameable only by its 'Ref' (opaque) or by
guessing the path back from its 'nodeShorthand' and hoping there is exactly
one node with that shorthand. A node reached by more than one seed, or twice
within one seed's graph, can have more than one path; all of them are shown,
since any one is a valid selector. Sourced from 'worldEpochs' only, same as
'resolveWorldSelectors' — a retired seed's graph is gone, and a node with no
entry here (nothing in it, or absent from the map) is one no /live/
declaration's graph currently reaches by path at all, addressable only by its
'Ref' (the @#@-prefixed form 'Query.resolveRewrittenSelectors' understands).
-}
worldPaths :: World seed directive -> Map Ref [Text]
worldPaths w =
Map.map (nub . sortOn Text.length) $
Map.fromListWith
(++)
[ (ref, [Text.intercalate "/" path])
| ep <- w.worldEpochs
, (path, ref) <- Query.pathedRefs ep.epochGraph
]