bus

package
v1.2.11 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 11, 2026 License: MIT Imports: 19 Imported by: 0

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

View Source
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).

View Source
const (
	LogKey = "bus2:log"
	Prefix = "bus2:to:"
)

The keys (SPEC-BUS.md, the data): one stream per recipient, and one log.

View Source
const (
	MaxBody = 1 << 20 // bytes of a body
	MaxName = 64      // bytes of a name
)

Limits (SPEC-BUS.md, the data).

View Source
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.

View Source
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.

View Source
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.

View Source
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).

View Source
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).

View Source
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).

View Source
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).

View Source
const MaxToken = 128

MaxToken is the bytes of a token at most.

View Source
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).

View Source
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.

View Source
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.

View Source
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.

View Source
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.

View Source
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.

View Source
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

Kinds is every kind, in the order the help lists them.

View Source
var StageStates = []string{Delivered, Read, Acted}

StageStates is every state, in order.

View Source
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

func CheckAddr(ctx context.Context, addr string, lookup Lookup) string

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

func CheckKinds(ks ...string) string

CheckKinds says why ks are no kinds, "" when each is one.

func CheckName

func CheckName(s string) string

CheckName says why s is no name, "" when it is one.

func CheckToken

func CheckToken(s string) string

CheckToken says why s is no token, "" when it is one; the empty token is a send with none.

func IDAt

func IDAt(t time.Time) string

IDAt is the first entry id a stream could hold at t (<ms>-0): the floor of a Range from that instant.

func OwedOf

func OwedOf(name string) string

OwedOf is the friend's hash of messages owed her session's receipt.

func SentOf

func SentOf(from, token string) string

SentOf is the key of the record of from's token.

func StagesOf

func StagesOf(name string) string

StagesOf is the recipient's receipts hash.

func StreamOf

func StreamOf(name string) string

StreamOf is the recipient's stream.

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.

func (Alarm) Text

func (a Alarm) Text() string

Text is the alarm as the coordinator reads it, one line.

type Backlog

type Backlog struct {
	Name     string
	Count    int
	OldestID string
	OldestAt time.Time
}

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.

func (Backlog) Age

func (bl Backlog) Age(now time.Time) time.Duration

Age is how long the oldest message has waited for its receipt at now; zero when nothing is owed.

type Bus

type Bus struct {
	Store Store
	// Rand fills a ULID's random half; crypto/rand when nil. Every Send calls
	// it, so a Bus shared by concurrent senders needs one safe for concurrent use.
	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

func (b *Bus) Ack(ctx context.Context, as string, ids []string) (map[string]bool, error)

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

func (b *Bus) AckEntry(ctx context.Context, as, entry string) (acked bool, err error)

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

func (b *Bus) Check(ctx context.Context, m Message) (Message, error)

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

func (b *Bus) Enroll(ctx context.Context, friends ...string) (added []string, err error)

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

func (b *Bus) Heard(ctx context.Context, names ...string) error

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

func (b *Bus) Log(ctx context.Context, from string) ([]Entry, error)

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

func (b *Bus) LogCursor(ctx context.Context) (string, error)

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) Names

func (b *Bus) Names(ctx context.Context) ([]string, error)

Names is the roster, sorted.

func (*Bus) Overdue

func (b *Bus) Overdue(ctx context.Context, older time.Duration) ([]Late, time.Time, error)

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

func (b *Bus) Peek(ctx context.Context, as string) (pending, fresh []Entry, err error)

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

func (b *Bus) ProvePush(ctx context.Context, p PushProof) (PushProof, error)

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

func (b *Bus) PushProofs(ctx context.Context, names ...string) ([]PushProof, time.Time, error)

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

func (b *Bus) Receipt(ctx context.Context, as string, ids []string) (int, error)

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

func (b *Bus) Send(ctx context.Context, m Message) (Message, error)

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

func (b *Bus) Stages(ctx context.Context, as string, ids ...string) ([]Stage, time.Time, error)

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

func (b *Bus) Stamp(ctx context.Context, as, state string, ids ...string) ([]string, error)

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

func (b *Bus) Undelivered(ctx context.Context, names ...string) ([]Backlog, error)

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

func (b *Bus) Unheard(ctx context.Context, names ...string) ([]string, error)

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

func (b *Bus) UnheardAt(ctx context.Context, at time.Time, names ...string) ([]string, error)

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

func (b *Bus) WaitArm(ctx context.Context, as, after string) (string, error)

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.

func (*Bus) Waiter

func (b *Bus) Waiter() (Waiter, error)

Waiter is the Store's wait reads, or the refusal of a Store that has none.

func (*Bus) WouldAck

func (b *Bus) WouldAck(ctx context.Context, as string, ids []string) (map[string]bool, error)

WouldAck is Ack that writes nothing (an ack's --dry-run): each id true when it is pending for the recipient, so Ack would ack it.

type Enroller

type Enroller interface {
	Enroll(ctx context.Context, friends ...string) error
}

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

func FilterKinds(es []Entry, kinds []string) []Entry

FilterKinds is the entries whose message is one of kinds; with no kinds, all of them.

func WaitPick

func WaitPick(entries []Entry, me string, skips []string) (kept []Entry, after string)

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.

func (Entry) Message

func (e Entry) Message() Message

Message is the entry's message.

type Late

type Late struct {
	Name, ID, From, Subject, State string
	At                             time.Time
	Age                            time.Duration
}

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

type Lookup func(ctx context.Context, host string) ([]netip.Addr, error)

Lookup resolves a host name to its addresses: net.DefaultResolver's LookupNetIP in the tool, a table in a test.

type Mark

type Mark struct {
	Key, Field, Value string
	Clear, Forward    bool
}

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

func Parse(fields map[string]string) Message

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.

func (Message) Fields

func (m Message) Fields() map[string]string

Fields is the message as the stream entry holds it.

func (Message) KindName

func (m Message) KindName() string

KindName is the message's kind, status when it has none (a message sent before kinds).

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) Age

func (p PushProof) Age(now time.Time) time.Duration

Age is how long ago the proof was written, at now; zero when there is none.

func (PushProof) AgeWord

func (p PushProof) AgeWord(now time.Time) string

AgeWord is Age as names and the refusal say it: "never" when there is none.

func (PushProof) Deaf

func (p PushProof) Deaf(now time.Time) string

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) State

func (p PushProof) State(now time.Time) string

State is the proof's push state at now; a zero proof is PushNone.

func (PushProof) Unheard

func (p PushProof) Unheard(now time.Time) string

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) Ack

func (r Redis) Ack(ctx context.Context, stream, group string, ids ...string) (n int64, err error)

func (Redis) AddAll

func (r Redis) AddAll(ctx context.Context, streams []string, fields map[string]string, marks ...Mark) error

func (Redis) AddOnce

func (r Redis) AddOnce(ctx context.Context, key, record string, keep time.Duration, streams []string, fields map[string]string, marks ...Mark) (prior string, found bool, err error)

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) Claim

func (r Redis) Claim(ctx context.Context, stream, group, consumer string, minIdle time.Duration, count int) ([]Entry, error)

func (Redis) Enroll

func (r Redis) Enroll(ctx context.Context, friends ...string) error

Enroll adds the friends to the set `friends`, in one trip.

func (Redis) EnsureGroup

func (r Redis) EnsureGroup(ctx context.Context, stream, group string) error

func (Redis) Forward

func (r Redis) Forward(ctx context.Context, key, state string, ids ...string) ([]string, error)

func (Redis) Get

func (r Redis) Get(ctx context.Context, stream string, ids []string) ([]Entry, error)

func (Redis) Group

func (r Redis) Group(ctx context.Context, stream, group string) (string, bool, error)

func (Redis) Marks

func (r Redis) Marks(ctx context.Context, keys ...string) ([]map[string]string, error)

func (Redis) Members

func (r Redis) Members(ctx context.Context) (friends, machines []string, now time.Time, err error)

func (Redis) Pending

func (r Redis) Pending(ctx context.Context, stream, group string, count int) ([]string, error)

func (Redis) Range

func (r Redis) Range(ctx context.Context, stream, from, to string, count int) ([]Entry, error)

func (Redis) Read

func (r Redis) Read(ctx context.Context, stream, group, consumer string, block time.Duration, count int) ([]Entry, error)

func (Redis) Release

func (r Redis) Release(ctx context.Context, stream, group string, ids ...string) error

func (Redis) Roster

func (r Redis) Roster(ctx context.Context) ([]string, time.Time, error)

func (Redis) Sent

func (r Redis) Sent(ctx context.Context, key string) (v string, ok bool, err error)

func (Redis) Tail

func (r Redis) Tail(ctx context.Context, stream string) (id string, ok bool, err error)

Tail is the stream's last entry id (XINFO STREAM's last-generated-id); a stream that is not there answers so, and the caller arms at "0-0" (SPEC-BUS.md, the verbs: wait). The read changes nothing, so a stall is tried once more (SPEC-BUS.md, the deadlines).

func (Redis) Unmark

func (r Redis) Unmark(ctx context.Context, key string, fields ...string) (n int64, err error)

type Refusal

type Refusal struct{ Problems []string }

Refusal is a reason a verb could not run as asked: the input, not the store.

func (*Refusal) Error

func (r *Refusal) Error() string

type Stage

type Stage struct {
	ID    string
	State string
	At    time.Time
}

Stage is one message's receipt: its state ("" none) and when the store reached it.

func (Stage) Age

func (s Stage) Age(now time.Time) time.Duration

Age is how long the receipt has stood at now; zero with none.

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.

func (*Watch) Observe

func (w *Watch) Observe(now time.Time, err error)

Observe is one send's result at now: nil is a success.

func (*Watch) Open

func (w *Watch) Open() bool

Open says whether an alarm is raised and not yet cleared.

Directories

Path Synopsis

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL