Documentation
¶
Overview ¶
Package bus is the message bus between AIs over Redis streams (docs/SPEC-BUS.md; the delivery machine is tla/Bus2.tla). A message is one stream entry, written to every recipient's stream and to the log in one transaction; a recipient reads its stream through a consumer group, so a message is pending from the moment it is delivered until it is acked, and a reader that crashes before acking is handed it again. The logic lives here, apart from the transport: Bus holds the rules over a Store, the few Redis commands the bus uses (redis.go is the Redis one, fake.go the one tests run on), so every rule is tested with no socket.
Index ¶
- Constants
- Variables
- func CheckAddr(ctx context.Context, addr string, lookup Lookup) string
- func CheckKinds(ks ...string) string
- func CheckName(s string) string
- func CheckToken(s string) string
- func IDAt(t time.Time) string
- func OwedOf(name string) string
- func SentOf(from, token string) string
- func StagesOf(name string) string
- func StreamOf(name string) string
- type Alarm
- type Backlog
- type Bus
- func (b *Bus) Ack(ctx context.Context, as string, ids []string) (map[string]bool, error)
- func (b *Bus) AckEntry(ctx context.Context, as, entry string) (acked bool, err error)
- func (b *Bus) Check(ctx context.Context, m Message) (Message, error)
- func (b *Bus) Enroll(ctx context.Context, friends ...string) (added []string, err error)
- func (b *Bus) Heard(ctx context.Context, names ...string) error
- func (b *Bus) Log(ctx context.Context, from string) ([]Entry, error)
- func (b *Bus) LogCursor(ctx context.Context) (string, error)
- func (b *Bus) LogForward(ctx context.Context, cursor string) (es []Entry, next string, more bool, err error)
- func (b *Bus) Names(ctx context.Context) ([]string, error)
- func (b *Bus) Overdue(ctx context.Context, older time.Duration) ([]Late, time.Time, error)
- func (b *Bus) Peek(ctx context.Context, as string) (pending, fresh []Entry, err error)
- func (b *Bus) ProvePush(ctx context.Context, p PushProof) (PushProof, error)
- func (b *Bus) PushProofs(ctx context.Context, names ...string) ([]PushProof, time.Time, error)
- func (b *Bus) Receipt(ctx context.Context, as string, ids []string) (int, error)
- func (b *Bus) Recv(ctx context.Context, as string, block time.Duration) (e Entry, ok bool, err error)
- func (b *Bus) RecvKinds(ctx context.Context, as string, block time.Duration, kinds []string) (e Entry, ok bool, err error)
- func (b *Bus) Send(ctx context.Context, m Message) (Message, error)
- func (b *Bus) Stages(ctx context.Context, as string, ids ...string) ([]Stage, time.Time, error)
- func (b *Bus) Stamp(ctx context.Context, as, state string, ids ...string) ([]string, error)
- func (b *Bus) Undelivered(ctx context.Context, names ...string) ([]Backlog, error)
- func (b *Bus) Unheard(ctx context.Context, names ...string) ([]string, error)
- func (b *Bus) UnheardAt(ctx context.Context, at time.Time, names ...string) ([]string, error)
- func (b *Bus) WaitArm(ctx context.Context, as, after string) (string, error)
- func (b *Bus) Waiter() (Waiter, error)
- func (b *Bus) WouldAck(ctx context.Context, as string, ids []string) (map[string]bool, error)
- type Enroller
- type Entry
- type Late
- type Lookup
- type Mark
- type Message
- type PushProof
- type Redis
- func (r Redis) Ack(ctx context.Context, stream, group string, ids ...string) (n int64, err error)
- func (r Redis) AddAll(ctx context.Context, streams []string, fields map[string]string, marks ...Mark) error
- func (r Redis) AddOnce(ctx context.Context, key, record string, keep time.Duration, streams []string, ...) (prior string, found bool, err error)
- func (r Redis) BlockRead(ctx context.Context, stream, after string, block time.Duration, count int) ([]Entry, error)
- func (r Redis) Claim(ctx context.Context, stream, group, consumer string, minIdle time.Duration, ...) ([]Entry, error)
- func (r Redis) Enroll(ctx context.Context, friends ...string) error
- func (r Redis) EnsureGroup(ctx context.Context, stream, group string) error
- func (r Redis) Forward(ctx context.Context, key, state string, ids ...string) ([]string, error)
- func (r Redis) Get(ctx context.Context, stream string, ids []string) ([]Entry, error)
- func (r Redis) Group(ctx context.Context, stream, group string) (string, bool, error)
- func (r Redis) Marks(ctx context.Context, keys ...string) ([]map[string]string, error)
- func (r Redis) Members(ctx context.Context) (friends, machines []string, now time.Time, err error)
- func (r Redis) Pending(ctx context.Context, stream, group string, count int) ([]string, error)
- func (r Redis) Range(ctx context.Context, stream, from, to string, count int) ([]Entry, error)
- func (r Redis) Read(ctx context.Context, stream, group, consumer string, block time.Duration, ...) ([]Entry, error)
- func (r Redis) Release(ctx context.Context, stream, group string, ids ...string) error
- func (r Redis) Roster(ctx context.Context) ([]string, time.Time, error)
- func (r Redis) Sent(ctx context.Context, key string) (v string, ok bool, err error)
- func (r Redis) Tail(ctx context.Context, stream string) (id string, ok bool, err error)
- func (r Redis) Unmark(ctx context.Context, key string, fields ...string) (n int64, err error)
- type Refusal
- type Stage
- type Store
- type TimeoutError
- type Waiter
- type Watch
Constants ¶
const ( AlarmAuth = "auth" AlarmConnection = "connection" )
The classes of a send alarm: the store refused the login, or could not be reached (SPEC-BUS.md, fr-delivery-receipts.w1).
const ( LogKey = "bus2:log" Prefix = "bus2:to:" )
The keys (SPEC-BUS.md, the data): one stream per recipient, and one log.
const ( MaxBody = 1 << 20 // bytes of a body MaxName = 64 // bytes of a name )
Limits (SPEC-BUS.md, the data).
const ( KindReport = "report" KindAck = "ack" KindStatus = "status" KindRequest = "request" KindBlocker = "blocker" )
The kinds of a message (SPEC-BUS.md, the kind of a message): the bus's own vocabulary, which a reader filters on; the bus gives none of them a meaning.
const ( PushProven = "proven" // up and younger than PushFresh: the name is heard PushStale = "stale" // up when written, and older than PushFresh: its daemon stopped renewing PushDown = "down" // its daemon says the session did not answer its check PushNone = "none" // no daemon ever recorded one )
The push states names shows, read off a proof at the store's now.
const ( CallTimeout = 5 * time.Second BlockMargin = 10 * time.Second )
The deadlines (SPEC-BUS.md, the deadlines). CallTimeout bounds every call that does not block; a blocking read gets the time it asked for and BlockMargin more, on top of Timeout. One retry, and only for a read that changes nothing.
const ( Delivered = "delivered" Read = "read" Acted = "acted" )
The states of a receipt, in the one order it moves: delivered (the recipient's reader took the message off its stream), read (the turn carrying it started), acted (the turn ended at exit 0, or the recipient sent a message naming it by re).
const ( // DefaultTokenLife is how long a retry under a token answers the // original: a day, longer than any sender's retry loop runs. DefaultTokenLife = 24 * time.Hour // DefaultTokenCleanup is when the store drops the record (its key // expires): a week. Between the life and the cleanup a retry is refused, // naming the message that went, never sent again. DefaultTokenCleanup = 7 * 24 * time.Hour )
The token's settings, when the Bus names none (Bus.TokenLife, TokenCleanup).
const ClaimAfter = 15 * time.Minute
ClaimAfter is how long a delivered message stays with its reader before recv hands it to another. It is longer than the longest delivery a reader makes (nova-friend's ten minute turn and the kill that ends it), so a live reader mid-turn is never handed its message a second time; a dead one's is claimed after this (SPEC-BUS.md, the semantics; tla/Bus2.tla HeldStaysHeld).
const Consumer = "nova-bus2"
Consumer is the one consumer name of every reader: with ClaimAfter, who holds an entry is told by its idle time, never by a name. It keeps the bus2 spelling with the keys: a consumer name lives in the live store's pending lists, and renaming it there is a migration (SPEC-BUS.md, the data).
const MaxToken = 128
MaxToken is the bytes of a token at most.
const OwedPrefix = "bus2:owed:"
OwedPrefix is the hash of what a friend is owed a receipt for: one field per message id, its value the message's at (RFC 3339), written in the send's own transaction (SPEC-BUS.md, fr-delivery-receipts.w1).
const PushFresh = 10 * time.Minute
PushFresh is how young a proof must be for its name to read as heard: a daemon that stopped renewing is a push nobody has proven for this long.
const PushKey = "bus2:push"
PushKey is the hash of every name's inbox push proof: one field per name, its value the PushProof as JSON (SPEC-BUS.md, bus-requires-inbox-push-proof). The proof is the seat's liveness rule, and advice to a sender: the friend daemon writes it when its session answers a SESSION CHECK carried in by the harness's deliver adapter, renews it while the session stays up, and writes it down when a check goes unanswered; names shows it, and send and recv say it as a NOTE beside a message that landed or was read, never as a refusal.
const SentPrefix = "bus2:sent:"
SentPrefix is the record of a send's token: one string key per sender and token, its value the record's JSON, expiring at the token's cleanup.
const StagePrefix = "bus2:receipt:"
StagePrefix is the hash of a recipient's receipts: one field per message id, its value the state and the store's time it was reached, in Unix seconds ("read 1791288000"); written only by Store.Forward (SPEC-BUS.md, message-receipts; tla/Bus2Receipts.tla), by the receipt rule the store's script keeps (redis.go, forwardLua) and each fake keeps beside it: a receipt moves only forward, and only delivered starts one.
const WaitMax = 5
WaitMax bounds the message lines one wait prints (SPEC-BUS.md, the verbs: wait); the entries past them stay for the next run.
const WaitRead = 100
WaitRead bounds the entries one blocking read of a wait takes: one batch holds the skipped and the counted of one burst, and the decision over it is one pure function (SPEC-BUS.md, the verbs: wait).
Variables ¶
var Kinds = []string{KindReport, KindAck, KindStatus, KindRequest, KindBlocker}
Kinds is every kind, in the order the help lists them.
var StageStates = []string{Delivered, Read, Acted}
StageStates is every state, in order.
var Tailnet = netip.MustParsePrefix("100.64.0.0/10")
Tailnet is the one network beyond loopback the bus reaches a store over: the tailnet is the boundary, with no ACL behind it (decided 2026-10-04), so an address outside it is refused before anything is dialled.
Functions ¶
func CheckAddr ¶
CheckAddr is the rule on where a store may be: "" when addr (host:port, or the absolute path of a Unix socket, which is this machine's) is on loopback or the tailnet, else the one line that names the rule and what the address is. A name is resolved and every address it has must pass; a name that does not resolve fails the rule too, since what it would dial is unknown. The shape of the address is redisconn's to refuse, not this.
func CheckKinds ¶
CheckKinds says why ks are no kinds, "" when each is one.
func CheckToken ¶
CheckToken says why s is no token, "" when it is one; the empty token is a send with none.
func IDAt ¶
IDAt is the first entry id a stream could hold at t (<ms>-0): the floor of a Range from that instant.
Types ¶
type Alarm ¶
type Alarm struct {
Store, User string
Class string // AlarmAuth or AlarmConnection
Reason string
Since time.Time
Failures int
Cleared time.Time // set on the alarm Clear is handed
}
Alarm is one outage of sends on a bus store: who logged in where, why it failed, since when and how many sends failed in it. It never holds a password: the reason is the store's refusal word, or the transport's error, which carries the address and never the secret.
type Backlog ¶
Backlog is what one friend is owed receipts for: how many messages, and the oldest of them (its id and when it was sent, by the store's clock); zero values when nothing is owed.
type Bus ¶
type Bus struct {
Store Store
// Rand fills a ULID's random half; crypto/rand when nil.
Rand func([]byte) (int, error)
// TokenLife is how long a retry under a send's token answers the
// original message (DefaultTokenLife when zero); TokenCleanup is when the
// store drops the token's record (DefaultTokenCleanup when zero, never
// before the life ends). token.go.
TokenLife, TokenCleanup time.Duration
// OnStampError hears a delivered receipt recv stamps that the store did
// not write; the recv goes on, and the message stays in Overdue until a
// later stamp lands. Nil drops it, the overdue alarm standing for it
// (stages.go).
OnStampError func(err error)
}
Bus is the rules over a Store.
func (*Bus) Ack ¶
Ack acks messages by their ids: each is looked up among the recipient's pending entries, so an id that is not pending (acked already, never delivered, or not this recipient's) is answered acked=false, never a failure: ack is idempotent. Ack by id is the session's verb (nova-bus ack), so it is also the session's receipt of every id it names, pending or not (Receipt): a daemon that acked the stream first takes nothing from it.
func (*Bus) AckEntry ¶
AckEntry acks one entry the recipient was handed (XACK); acking it again is a no-op that says so. (tla/Bus2.tla: Ack, AckIdempotent)
func (*Bus) Check ¶
Check is Send that writes nothing (a send's --dry-run): every problem of the message named at once, as Send names them, and the message as it would be sent, at the store's time with its recipients sorted, and no id: an id is made for a message sent.
func (*Bus) Enroll ¶
Enroll makes each of friends a known name of the bus, a friend owed a receipt: the ones the roster does not hold are added, and added is them, sorted. A name that is no name (CheckName) is refused and the rest are added all the same. A store that cannot be told (a test's fake) is told nothing and adds none. Two trips: the roster, then the add when one is missing.
func (*Bus) Heard ¶
Heard refuses every name of names (deduplicated, in order) that is not heard at the store's now, one deaf line each: the gate nova-bus's --require-push puts in front of a send or a recv, before anything is written or read. It is the exception, asked for by flag; the bus itself never refuses on a proof.
func (*Bus) Log ¶
Log is the log's messages from the entry id from ("-" for its start), oldest first, up to logLimit of them. A reader that always passes "-" sees only the oldest window, so once the log is longer than logLimit a fresh entry is past the cap on every read. Arm at LogCursor and read with LogForward.
func (*Bus) LogCursor ¶
LogCursor is the log's tail, the cursor a reader arms at so LogForward returns only entries appended after it. An empty log arms at "0-0", the id before any entry. The store's Tail is the same read a wait arms with.
func (*Bus) LogForward ¶
func (b *Bus) LogForward(ctx context.Context, cursor string) (es []Entry, next string, more bool, err error)
LogForward is one bounded window of the log strictly after cursor, oldest first, at most logLimit entries. cursor is an entry id ("0-0" before any). next is the last entry read, or cursor when the window is empty, so the caller advances as it consumes and a fresh entry is not hidden behind the oldest logLimit. more is set when the window is full and a further read from next may hold more.
func (*Bus) Overdue ¶
Overdue is every message on every known name's stream still short of delivered older than older, oldest first, and the store's time: the alarm the coordinator's loop runs. A message on a stream is short of delivered when it is new, or pending with no receipt (its reader's stamp did not land); one acked is never listed. One trip for the roster, a peek of each stream, and one for every receipts hash.
func (*Bus) Peek ¶
Peek is what waits for the recipient, reading only: the pending entries (delivered, not acked) and the new ones (never delivered), oldest first, up to pendingLimit of each.
func (*Bus) ProvePush ¶
ProvePush records proof for name at the store's time, in one transaction (HSET on PushKey through AddAll, with no stream): what the friend daemon calls when its session answered a check carried in by the deliver adapter, while it stays up, and with Up false when a check went unanswered.
func (*Bus) PushProofs ¶
PushProofs is each name's proof, in the order asked, and the store's now they are read at (a roster trip, then one HGETALL).
func (*Bus) Receipt ¶
Receipt is the session's word that it read the messages ids sent to as: each is cleared from what as is owed, and the answer is how many were owed. An id not owed (received already, or never as's) is no failure: a receipt is idempotent. Only the session gives one (nova-bus ack by id, or a reply naming the message); the daemon's ack of the stream at the end of a turn is no receipt, so a turn the session never read stays undelivered.
func (*Bus) Recv ¶
func (b *Bus) Recv(ctx context.Context, as string, block time.Duration) (e Entry, ok bool, err error)
Recv is one message for the recipient: the oldest one delivered and not acked whose reader has had it longer than ClaimAfter (a reader that died or stalled), else the oldest never delivered, waiting up to block for it. ok is false when there is none. A name the roster does not hold is refused, never given a stream to wait on. The group is made on first use. (tla/Bus2.tla: Recv, PendingBeforeNew, HeldStaysHeld)
func (*Bus) RecvKinds ¶
func (b *Bus) RecvKinds(ctx context.Context, as string, block time.Duration, kinds []string) (e Entry, ok bool, err error)
RecvKinds is Recv for the messages whose kind is one of kinds (none: any). The message it hands out is stamped delivered on the recipient's receipts, one trip more, and carries the receipt it found (Entry.Stage): acted on one the claim hands in again after a turn acted on it (stages.go). A message the filter skips is handed back to the group at once (Store.Release), neither acked nor held, so a reader that asks for all gets it next; the skip costs one round trip per message skipped, and one more to release a run of them. (tla/Bus2.tla: Recv; a skipped message is back as a lost one)
func (*Bus) Send ¶
Send checks the message, stamps it with the store's time and a ULID made from that time, and appends it to every recipient's stream and the log in one transaction. It answers the message as sent. Every problem of the message is named at once in one Refusal: a name that is no name, a recipient the roster does not know (with how to add one), an empty body, a body over MaxBody, a from that is unknown.
A message to a friend is owed her session's receipt (receipt.go): the transaction marks it on bus2:owed:<friend> for each friend it names but the sender, and a message from a friend naming another (re) is her receipt of that one, cleared in the same transaction.
A message naming another (re) is the sender's act on it: its receipt on bus2:receipt:<sender> moves to acted in the same transaction, when the sender was delivered it (stages.go).
A message with a Token is sent once under it (token.go): the record of the token is written in the same step, and a send that finds it writes nothing and answers the message it names, id and at, so a caller whose response was lost after the write committed retries with the same token and the same arguments and gets the original. The same token with other arguments, or past its life, is refused. Without a token a lost response retried is a second message.
func (*Bus) Stages ¶
Stages is as's receipts and the store's time, in two trips: those of ids in their order (State "" for one with none), or with no ids every receipt as holds, oldest first.
func (*Bus) Stamp ¶
Stamp moves the receipts of ids on as's hash to state, each only forward (Forward), at the store's time, in one trip, and answers each id's state before it ("" none). It is the daemon's word that a turn started (read) or ended at exit 0 (acted); recv stamps delivered, and a send naming a message (re) stamps it acted in its own transaction (owe).
func (*Bus) Undelivered ¶
Undelivered is each name's Backlog, in the order asked, in one trip: what the friends table shows as undelivered and the oldest undelivered age. An at that is no instant still counts, and is oldest only when none is.
func (*Bus) Unheard ¶
Unheard is one advisory line for every name of names (deduplicated, in order) that is not heard at the store's now, in one trip after the roster's (one HGETALL): what send prints as SEND NOTE after its message landed and recv as RECV NOTE beside what it read. It is advice, never a gate: the push proof is the seat's liveness rule (names, the coordinator's view), and a message is never refused on it (the finding of 2026-10-08, issue #5450: a claude friend and a machine with no daemon could never be written to or read as).
func (*Bus) UnheardAt ¶
UnheardAt is Unheard judged at at, with no roster trip: send's, judged at the store's time its message was stamped with (one HGETALL after the write).
func (*Bus) WaitArm ¶
WaitArm is the cursor a wait on as starts from: after when given (the stream's tail is not read), else the stream's last entry id read once ("0-0" when the stream is not there). A name the roster does not hold is refused, never given a stream to wait on (SPEC-BUS.md, the semantics), as recv refuses one.
type Enroller ¶
Enroller is a Store whose roster can be told friends: their names added to the set `friends` (SADD), never one taken off. nova-config's apply writes the friend rows to the sprint store, and the bus store may be another Redis (the fleet runs it apart, NOVA_BUS_REDIS beside NOVA_SPRINT_REDIS), so the one who reads the rows (nova-sprint friend sync) tells the bus store (SPEC-BUS.md, the config). A row removed stays a known name here: the bus never deletes.
type Entry ¶
type Entry struct {
Stream string
Entry string
Fields map[string]string
// Stage is the message's receipt state as recv found it, before its
// delivered stamp ("" none; stages.go): acted on a message the claim hands
// in again after a turn acted on it. Only recv sets it.
Stage string
}
Entry is one stream entry as read: the stream, the entry's id in it (<ms>-<seq>, the server's) and its fields.
func FilterKinds ¶
FilterKinds is the entries whose message is one of kinds; with no kinds, all of them.
func WaitPick ¶
WaitPick is the wait's decision over one batch of entries, a pure function (SPEC-BUS.md, the verbs: wait): the first WaitMax entries not from me whose subject starts with none of skips (matched without case), in order, and the cursor past every entry the walk saw -- a skipped entry moves it -- so a caller that re-arms with it misses nothing between runs. The walk stops at the WaitMax-th entry that counts, and the entries after it stay for the next run.
type Late ¶
Late is a message still short of delivered: on name's stream, new (never delivered) or pending with no receipt, and how long since it was sent.
type Lookup ¶
Lookup resolves a host name to its addresses: net.DefaultResolver's LookupNetIP in the tool, a table in a test.
type Mark ¶
Mark is one write to a hash inside a send's transaction: HSET Key Field Value, or HDEL Key Field when Clear, or, when Forward, the receipt of Field moved to the state Value by the receipt rule (redis.go, forwardLua).
type Message ¶
type Message struct {
ID string
From string
To []string
CC []string
Subject string
Re string
Kind string // one of Kinds; "" is status
At time.Time
Body string
// Token is the caller's word for this one logical send, the same on
// every retry of it (token.go); it is the sender's, never on the entry.
// Empty is a send with none: every call a new message.
Token string
}
Message is one message of the bus: its fields on every stream it is on.
func Parse ¶
Parse is the message an entry's fields hold. A field that is not there is empty (a kind that is not there is read as status by KindName); an `at` that is no instant is the zero time, never a refusal, so a log with one odd entry still reads.
type PushProof ¶
type PushProof struct {
Name string `json:"-"`
Harness string `json:"harness"`
Nonce string `json:"nonce"`
Proven time.Time `json:"proven"`
Up bool `json:"up"`
Reason string `json:"reason,omitempty"`
At time.Time `json:"at"`
}
PushProof is one name's proof that its inbox pushes into its session: the harness whose deliver adapter carried the check, the nonce the session answered and when (Proven, the daemon's clock), whether the daemon still holds the session up (Up; Reason when it does not), and when the daemon last wrote it (At, the store's clock, the one freshness is read by).
func (PushProof) AgeWord ¶
AgeWord is Age as names and the refusal say it: "never" when there is none.
func (PushProof) Deaf ¶
Deaf is why name is not heard at now, "" when its proof is proven: the refusal of nova-bus's --require-push, with the remedy.
func (PushProof) Unheard ¶
Unheard is the advisory line for name at now, "" when its proof is proven: the NOTE send and recv print beside their result (SPEC-BUS.md, bus-requires-inbox-push-proof). It never refuses: a message to an unheard name lands and waits on its stream; the line says what the proof's state is and what would prove one, so a sender knows nothing is pushing it in.
type Redis ¶
type Redis struct {
C *redis.Client
Timeout time.Duration
Margin time.Duration // a blocking read's wait beyond its block (BlockMargin when zero)
}
Redis is the Store over one go-redis client. Each method is one round trip, under a deadline of Timeout (CallTimeout when zero); the client must honour a context's deadline (redisconn sets ContextTimeoutEnabled).
func (Redis) BlockRead ¶
func (r Redis) BlockRead(ctx context.Context, stream, after string, block time.Duration, count int) ([]Entry, error)
BlockRead is XREAD past the id after, waiting up to block (0 is for ever: one read the store holds until an entry is there), never the consumer group, so a later recv still delivers and acks what it handed out (SPEC-BUS.md, the verbs: wait). A positive block is bounded by Timeout plus the block plus the margin, and is not retried; block 0 keeps the caller's context, because a bus deadline would end a wait that asked to park (SPEC-BUS.md, the deadlines).
func (Redis) EnsureGroup ¶
type Refusal ¶
type Refusal struct{ Problems []string }
Refusal is a reason a verb could not run as asked: the input, not the store.
type Store ¶
type Store interface {
// Roster is the known names (nova-config's friend and machine rows, the
// sets `friends` and `machines`) and the server's time (TIME), in one trip.
Roster(ctx context.Context) (names []string, now time.Time, err error)
// Members is the friends and the machines apart (the sets `friends` and
// `machines`) and the server's time (TIME), in one trip: Roster split, so
// a send knows which recipients are friends, owed a receipt.
Members(ctx context.Context) (friends, machines []string, now time.Time, err error)
// AddAll appends one entry with fields to every stream, and makes every
// mark (HSET, or HDEL when it clears), in one MULTI/EXEC: the entry and its
// marks are on all of them or on none.
AddAll(ctx context.Context, streams []string, fields map[string]string, marks ...Mark) error
// AddOnce is AddAll under a send's token, in one atomic step (a script):
// when key holds a record it writes nothing and answers that record and
// found; else it sets key to record, expiring after keep, and appends
// the entry and makes the marks as AddAll does. The record's key is the
// one key the bus writes that the store removes, by its expiry.
AddOnce(ctx context.Context, key, record string, keep time.Duration, streams []string, fields map[string]string, marks ...Mark) (prior string, found bool, err error)
// Sent is the record at key (GET), and whether there is one.
Sent(ctx context.Context, key string) (record string, found bool, err error)
// Unmark clears fields of the hash at key (HDEL) and says how many were there.
Unmark(ctx context.Context, key string, fields ...string) (int64, error)
// Marks is the whole hash at each key, in one trip (a pipeline of HGETALL);
// a key that is not there is an empty map.
Marks(ctx context.Context, keys ...string) ([]map[string]string, error)
// EnsureGroup makes the group on the stream from its start, making the
// stream when it is not there (XGROUP CREATE ... 0 MKSTREAM); a group
// already there is fine.
EnsureGroup(ctx context.Context, stream, group string) error
// Claim hands consumer up to count entries pending for the group that have
// been idle (delivered and not acked) for at least minIdle (XAUTOCLAIM
// <minIdle> 0-0): what a reader that died, or stalled, was holding.
Claim(ctx context.Context, stream, group, consumer string, minIdle time.Duration, count int) ([]Entry, error)
// Read hands consumer up to count entries the group has never delivered
// (XREADGROUP ... >), waiting up to block for one when block is above
// zero, else answering at once.
Read(ctx context.Context, stream, group, consumer string, block time.Duration, count int) ([]Entry, error)
// Release makes the entries pending for the group claimable at once
// (XCLAIM ... IDLE <ClaimAfter> JUSTID): what a reader that skipped them
// hands back, in one trip.
Release(ctx context.Context, stream, group string, entries ...string) error
// Ack acks entries for the group (XACK) and says how many were pending.
Ack(ctx context.Context, stream, group string, entries ...string) (int64, error)
// Pending is the entry ids pending for the group, up to count (XPENDING).
Pending(ctx context.Context, stream, group string, count int) ([]string, error)
// Group is the group's last delivered entry id, and whether the group is
// there at all (XINFO GROUPS).
Group(ctx context.Context, stream, group string) (lastDelivered string, exists bool, err error)
// Range is the entries of the stream from from to to, up to count
// (XRANGE; "-", "+" and the exclusive "(<id>" as Redis reads them).
Range(ctx context.Context, stream, from, to string, count int) ([]Entry, error)
// Get is the named entries of the stream, in one trip (a pipeline of XRANGE
// id id); an id that is not there is left out.
Get(ctx context.Context, stream string, entries []string) ([]Entry, error)
// Forward moves each id's receipt on the hash at key to state at the
// store's time (TIME), by the receipt rule (forwardLua): only forward, and only
// delivered starts one; it answers each id's state before, "" for none, in
// one atomic step (a script). It is the one writer of a receipt.
Forward(ctx context.Context, key, state string, ids ...string) ([]string, error)
}
Store is the few Redis commands the bus uses, each one round trip. The bus never deletes: no command here removes an entry, a group or a key (a token's record expires: AddOnce).
type TimeoutError ¶
type TimeoutError struct {
Addr string
After time.Duration
// contains filtered or unexported fields
}
TimeoutError is a call the store did not answer within After. It names the address and never the login.
func (*TimeoutError) Error ¶
func (e *TimeoutError) Error() string
func (*TimeoutError) Temporary ¶
func (e *TimeoutError) Temporary() bool
func (*TimeoutError) Timeout ¶
func (e *TimeoutError) Timeout() bool
func (*TimeoutError) Unwrap ¶
func (e *TimeoutError) Unwrap() error
type Waiter ¶
type Waiter interface {
// Tail is the stream's last entry id and whether the stream is there at
// all (XINFO STREAM's last-generated-id, "0-0" for an empty one): the
// cursor a wait arms at when the caller gives none (SPEC-BUS.md, the
// verbs: wait).
Tail(ctx context.Context, stream string) (last string, exists bool, err error)
// BlockRead hands up to count entries of the stream lying past the id
// after (XREAD), waiting up to block for one when block is above zero (0
// is for ever), else answering at once. It never touches the consumer
// group, so what it hands out is still a later recv's to deliver and ack
// (SPEC-BUS.md, the verbs: wait).
BlockRead(ctx context.Context, stream, after string, block time.Duration, count int) ([]Entry, error)
}
Waiter is the two reads a wait makes over a Store that also holds them: the Redis store does; a Store without them cannot wait (SPEC-BUS.md, the verbs: wait).
type Watch ¶
type Watch struct {
Store, User string
Raise, Clear func(Alarm)
// contains filtered or unexported fields
}
Watch turns the results of sends on one store into alarms: the first send that fails on the login or the connection raises one alarm, every later failure only counts in it, and the next send that succeeds clears it. A refusal of the message (Refusal) or any other answer of the store is the sender's to read, neither a failure nor a success here. Raise and Clear are called with the Watch's lock held and must not call back into it.