Documentation
¶
Index ¶
- func MarshalGatePreparedRecord(rec GatePreparedRecord) ([]byte, error)
- func MarshalLeaseFence(f LeaseFence) ([]byte, error)
- func ValidateCommandRecordRoute(record CommandRecord) error
- type AmbiguousAckError
- type AppendError
- type AppendFunc
- type AppendMiddleware
- type AppendResult
- type AppenderOption
- type CommandRecord
- func (r CommandRecord) Command() command.Command
- func (r CommandRecord) DeliveryPhase() command.DelegateDeliveryPhase
- func (r CommandRecord) IdempotencyID() string
- func (r CommandRecord) LogicalCommandID() uuid.UUID
- func (r CommandRecord) LoopID() uuid.UUID
- func (r CommandRecord) NormalizedDeliveryFingerprint() (Fingerprint, error)
- func (r CommandRecord) PhysicalID() CommandRecordID
- func (r CommandRecord) SessionID() uuid.UUID
- type CommandRecordID
- type CommandRouteMismatchError
- type DeliveryTransitionError
- type EventCursor
- type EventRecord
- type EventReplayer
- type FenceDecodeError
- type FenceEncodeError
- type FenceRecord
- type Fingerprint
- type FollowUnsupportedError
- type GatePreparedDecodeError
- type GatePreparedEncodeError
- type GatePreparedRecord
- type IdempotencyCollisionError
- type IdempotencyIndex
- type IdempotentJournal
- type JournalCommandAppender
- type JournalEventAppender
- type JournalGateAppender
- type JournalLeaseLostError
- type JournalNotReadyError
- type JournalRecord
- type Lease
- type LeaseFence
- type LeaseHeldError
- type LeaseLostError
- type MarshalRecordError
- type NilJournalError
- type RecordCursor
- type RecordKindError
- type RecordReplayer
- type RecordTooLargeError
- type ReplayRequest
- type SessionJournal
- type StartPos
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
func MarshalGatePreparedRecord ¶
func MarshalGatePreparedRecord(rec GatePreparedRecord) ([]byte, error)
MarshalGatePreparedRecord encodes a GatePreparedRecord into the JSON body that the sessionstore envelope carries: the GatePrepared event (via event.MarshalEvent) and the sealed gate.OpenPayload (via gate.MarshalPayload), both as raw sibling keys. A nil payload is rejected fail-closed — a prepared record without its validation payload is corrupt.
func MarshalLeaseFence ¶
func MarshalLeaseFence(f LeaseFence) ([]byte, error)
MarshalLeaseFence encodes a LeaseFence as its minimal JSON object {"epoch":N}.
func ValidateCommandRecordRoute ¶
func ValidateCommandRecordRoute(record CommandRecord) error
ValidateCommandRecordRoute validates the command's own identity contract, then cross-checks the duplicated live dispatch route for machine NoFold or phased delegate input. A zero record LoopID is accepted only because storage replay cannot reconstruct it.
Types ¶
type AmbiguousAckError ¶
AmbiguousAckError reports an append whose outcome the backend could not resolve: the persist call was lost or timed out and a bounded retry stayed ambiguous, so the serializer cannot tell whether the record landed. The fence stays unadvanced, so the next Append re-fences on the same tip; the caller decides whether to fail the session or retry later. It carries the record's routing/destination identifier, its idempotency id, the expected sequence the append fenced on, and the underlying cause.
func (*AmbiguousAckError) Error ¶
func (e *AmbiguousAckError) Error() string
func (*AmbiguousAckError) Unwrap ¶
func (e *AmbiguousAckError) Unwrap() error
type AppendError ¶
AppendError wraps a definite failure to persist a record to the session log. It carries the record's routing/destination identifier, its idempotency id, and the expected sequence the append was fenced under, and unwraps to the underlying backend error (a context deadline, a fence rejection, a transport error). The fence stays unadvanced when this is returned, so the next Append re-fences on the same tip.
func (*AppendError) Error ¶
func (e *AppendError) Error() string
func (*AppendError) Unwrap ¶
func (e *AppendError) Unwrap() error
type AppendFunc ¶
type AppendFunc func(context.Context, JournalRecord) (uint64, error)
AppendFunc is one journal append operation.
type AppendMiddleware ¶
type AppendMiddleware func(next AppendFunc) AppendFunc
AppendMiddleware decorates one AppendFunc. Implementations must invoke next synchronously exactly once with the supplied record and return its exact sequence and error.
func HookMiddleware ¶
func HookMiddleware(runner *hook.Runner, sessionID uuid.UUID) AppendMiddleware
HookMiddleware observes safe, classifiable journal records with runner. Records whose metadata cannot be derived without panicking bypass observation and delegate unchanged.
type AppendResult ¶
AppendResult reports the outcome of an append issued through an IdempotentJournal: the durable sequence the record occupies, and whether this call durably persisted a NEW frame (Appended=true) or deduplicated an identical retry of an already-durable record (Appended=false; Sequence is then the ORIGINAL append's sequence, not a new one).
type AppenderOption ¶
type AppenderOption func(*JournalEventAppender)
AppenderOption configures a JournalEventAppender at construction. Applied in order over a defaults struct (nop catalog), so a later option overrides an earlier one.
func WithCatalog ¶
func WithCatalog(c catalogUpdater) AppenderOption
WithCatalog injects the catalog updater the appender notifies after a successful append (best-effort). A nil updater is ignored (the nop default is kept), so the appender owns its invariant — it never holds a nil catalog and never nil-derefs.
type CommandRecord ¶
type CommandRecord struct {
// contains filtered or unexported fields
}
CommandRecord wraps a command.Command targeting a specific loop. Unlike an event, a command does not uniformly carry its routing coordinates (Interrupt and Shutdown carry only a Header; the session dispatches them by other means), so the writer supplies the target sessionID/loopID at construction. Its id is the command's physical CommandRecordID. The serializer encodes the wrapped command via command.MarshalCommand.
func NewCommandRecord ¶
func NewCommandRecord(sessionID, loopID uuid.UUID, cmd command.Command) CommandRecord
NewCommandRecord wraps cmd as the intent-log record targeting loop loopID in session sessionID. The caller is the writer, which knows the dispatch target; the command itself may not carry it.
func (CommandRecord) Command ¶
func (r CommandRecord) Command() command.Command
Command returns the wrapped command for the serializer to marshal.
func (CommandRecord) DeliveryPhase ¶
func (r CommandRecord) DeliveryPhase() command.DelegateDeliveryPhase
DeliveryPhase returns the durable delegate-delivery phase carried by a UserInput, or the zero phase for every other command kind.
func (CommandRecord) IdempotencyID ¶
func (r CommandRecord) IdempotencyID() string
IdempotencyID is the command's physical id rendered canonically.
func (CommandRecord) LogicalCommandID ¶
func (r CommandRecord) LogicalCommandID() uuid.UUID
LogicalCommandID returns the command UUID independent of its durable phase.
func (CommandRecord) LoopID ¶
func (r CommandRecord) LoopID() uuid.UUID
LoopID is the dispatch target the writer recorded for this command — the loop the intent-log entry belongs to. It is the backend-neutral routing coordinate a consumer keys on (replacing the subject a NATS backend derived it into).
func (CommandRecord) NormalizedDeliveryFingerprint ¶
func (r CommandRecord) NormalizedDeliveryFingerprint() (Fingerprint, error)
NormalizedDeliveryFingerprint fingerprints the exact UserInput payload with only DelegateDeliveryPhase cleared. Accepted is json:"-" and therefore does not enter the fingerprint; blocks, route, agency, hand-back, timestamps, and all other durable fields remain part of it.
func (CommandRecord) PhysicalID ¶
func (r CommandRecord) PhysicalID() CommandRecordID
PhysicalID returns the typed idempotency identity used by the journal envelope and backend message id. Keeping phase selection here prevents callers from constructing unsafe ad-hoc string ids.
func (CommandRecord) SessionID ¶
func (r CommandRecord) SessionID() uuid.UUID
SessionID is the session this command was recorded under.
type CommandRecordID ¶
type CommandRecordID struct {
CommandID uuid.UUID
Phase command.DelegateDeliveryPhase
}
CommandRecordID is the physical idempotency identity of one command record. Delivery fallback records retain the logical command UUID while using their phase as a typed suffix, so the intent and fallback frames can both be durable without weakening the journal's ordinary idempotency collision rule.
func (CommandRecordID) String ¶
func (id CommandRecordID) String() string
String returns the canonical physical id. Ordinary and intent commands use the logical UUID unchanged; only fallback_queued has a distinct physical id.
type CommandRouteMismatchError ¶
CommandRouteMismatchError reports disagreement between a durable delegate command's embedded target and the live CommandRecord dispatch route.
func (*CommandRouteMismatchError) Error ¶
func (e *CommandRouteMismatchError) Error() string
type DeliveryTransitionError ¶
type DeliveryTransitionError struct {
CommandID uuid.UUID
Phase command.DelegateDeliveryPhase
Reason string
}
DeliveryTransitionError reports a phased delegate command that violates the logical intent-to-fallback ordering or otherwise cannot participate in the transition index. It carries only bounded identity/category data; payloads are deliberately excluded.
func (*DeliveryTransitionError) Error ¶
func (e *DeliveryTransitionError) Error() string
type EventCursor ¶
type EventCursor interface {
// Next returns the next event and its sequence, or io.EOF when the cold backlog is
// exhausted. A decode/read error fails secure: the cursor surfaces the typed error
// rather than skipping or zero-valuing the record.
Next(ctx context.Context) (event.Event, uint64, error)
// Close tears down the reader. Idempotent: a second call is a no-op.
Close() error
}
EventCursor yields a session's Enduring events in sequence order. Next returns the next decoded event with its sequence, io.EOF once the backlog is drained (cold mode), or a typed error on a malformed/missing/corrupt record. Close releases the underlying reader; it is idempotent and safe to call after an error.
type EventRecord ¶
type EventRecord struct {
// contains filtered or unexported fields
}
EventRecord wraps an Enduring event.Event as a JournalRecord. The event already carries its producer Coordinates and Scope; its id is the event's EventID. The serializer encodes the wrapped event via event.MarshalEvent (which fails closed on an Ephemeral event); the record never re-encodes.
func NewEventRecord ¶
func NewEventRecord(ev event.Event) EventRecord
NewEventRecord wraps ev for the journal. ev must be an Enduring event; the Ephemeral check is the serializer's (event.MarshalEvent), not this wrapper's.
func (EventRecord) Event ¶
func (r EventRecord) Event() event.Event
Event returns the wrapped event for the serializer to marshal.
func (EventRecord) IdempotencyID ¶
func (r EventRecord) IdempotencyID() string
IdempotencyID is the event's EventID rendered canonically.
type EventReplayer ¶
type EventReplayer interface {
// Open binds a cursor over the session's events selected by req and positioned at
// req.From.
Open(ctx context.Context, req ReplayRequest) (EventCursor, error)
}
EventReplayer is the journal's read side: it opens an ordered cursor over a session's Enduring events. It is the narrow counterpart to SessionJournal (the write side) — a caller that only reads history depends on Open alone. The concrete implementation lives in a backend package (e.g. pkg/sessionstore over storage).
type FenceDecodeError ¶
FenceDecodeError wraps a failure to decode LeaseFence bytes at the untrusted restore boundary: malformed JSON, a wrong field type, a non-object, or trailing data after the object. The codec fails closed with this typed error so callers inspect the cause via errors.As rather than guessing an epoch.
func (*FenceDecodeError) Error ¶
func (e *FenceDecodeError) Error() string
func (*FenceDecodeError) Unwrap ¶
func (e *FenceDecodeError) Unwrap() error
type FenceEncodeError ¶
type FenceEncodeError struct{ Cause error }
FenceEncodeError wraps a failure to marshal a LeaseFence to JSON. A LeaseFence is a single uint64, so this is effectively unreachable, but the codec returns a typed error rather than dropping the json.Marshal error to satisfy the errors-are-typed contract at the package API.
func (*FenceEncodeError) Error ¶
func (e *FenceEncodeError) Error() string
func (*FenceEncodeError) Unwrap ¶
func (e *FenceEncodeError) Unwrap() error
type FenceRecord ¶
type FenceRecord struct {
// contains filtered or unexported fields
}
FenceRecord wraps a LeaseFence as a JournalRecord. The fence carries no session id of its own, so the writer supplies the target sessionID at construction; its idempotency id is the epoch.
func NewFenceRecord ¶
func NewFenceRecord(sessionID uuid.UUID, fence LeaseFence) FenceRecord
NewFenceRecord wraps fence as the fence record for session sessionID.
func (FenceRecord) Fence ¶
func (r FenceRecord) Fence() LeaseFence
Fence returns the wrapped LeaseFence for the serializer to marshal.
func (FenceRecord) IdempotencyID ¶
func (r FenceRecord) IdempotencyID() string
IdempotencyID is the epoch rendered as a decimal string.
func (FenceRecord) SessionID ¶
func (r FenceRecord) SessionID() uuid.UUID
SessionID is the session this fence marks a lease handover for.
type Fingerprint ¶
type Fingerprint struct {
// contains filtered or unexported fields
}
Fingerprint identifies a record's persisted kind and payload bytes — exactly what a backend durably writes — independent of any transient in-memory routing a record wrapper additionally carries (CommandRecord's session/loop dispatch target is never persisted, so it never enters a Fingerprint: a backend derives Fingerprint from the SAME (kind, codec-marshaled body) pair it writes to the log, never from the record's Go value). Two records fingerprint equal if and only if a backend would durably persist byte-identical frames for them.
func NewFingerprint ¶
func NewFingerprint(kind string, body []byte) Fingerprint
NewFingerprint derives the Fingerprint of a record's persisted envelope kind (the backend-neutral kind name, e.g. "event"/"command"/"fence") and its codec-marshaled payload bytes, exactly as a backend encodes them before persisting.
type FollowUnsupportedError ¶
type FollowUnsupportedError struct {
Stream string
}
FollowUnsupportedError is returned by Open when ReplayRequest.Follow is true and the backend implements only the cold (Follow:false) backlog read. Failing closed with a typed error is preferable to silently behaving as a cold cursor that EOFs at the current tip. It carries the log identifier the replay targeted.
func (*FollowUnsupportedError) Error ¶
func (e *FollowUnsupportedError) Error() string
type GatePreparedDecodeError ¶
type GatePreparedDecodeError struct {
Stage string // "json", "prepared", or "payload"
Cause error
}
GatePreparedDecodeError wraps a failure to decode GatePreparedRecord bytes at the untrusted restore boundary: malformed JSON, a malformed embedded event, or a malformed embedded payload. It fails secure rather than skipping or zero-valuing the record.
func (*GatePreparedDecodeError) Error ¶
func (e *GatePreparedDecodeError) Error() string
func (*GatePreparedDecodeError) Unwrap ¶
func (e *GatePreparedDecodeError) Unwrap() error
type GatePreparedEncodeError ¶
GatePreparedEncodeError wraps a failure to marshal a GatePreparedRecord to JSON: either the embedded GatePrepared event or the gate.OpenPayload failed its codec.
func (*GatePreparedEncodeError) Error ¶
func (e *GatePreparedEncodeError) Error() string
func (*GatePreparedEncodeError) Unwrap ¶
func (e *GatePreparedEncodeError) Unwrap() error
type GatePreparedRecord ¶
type GatePreparedRecord struct {
// contains filtered or unexported fields
}
GatePreparedRecord is the PRIVATE durable record for a gate's prepare step. It carries the GatePrepared event (the public envelope stored privately until ActivateGate appends the public GateOpened) PLUS the sealed gate.Payload the resolver needs for response validation and restore — a payload that must NEVER be exposed to SSE/history and must NEVER be appended through NewEventRecord or hub.PublishEvent. Its idempotency id is the prepared event's EventID.
func NewGatePreparedRecord ¶
func NewGatePreparedRecord(prepared event.GatePrepared, payload gate.OpenPayload) GatePreparedRecord
NewGatePreparedRecord wraps the private prepared projection and its typed payload as a single private journal record. The caller must NOT also append the GatePrepared event as a public EventRecord.
func UnmarshalGatePreparedRecord ¶
func UnmarshalGatePreparedRecord(data []byte) (GatePreparedRecord, error)
UnmarshalGatePreparedRecord decodes bytes produced by MarshalGatePreparedRecord back into a GatePreparedRecord. It fails closed with a typed *GatePreparedDecodeError on any malformed input — malformed JSON, a malformed embedded event, or a malformed embedded payload — so restore never silently drops or zero-values a private prepared record.
func (GatePreparedRecord) IdempotencyID ¶
func (r GatePreparedRecord) IdempotencyID() string
IdempotencyID is the prepared event's EventID rendered canonically.
func (GatePreparedRecord) Payload ¶
func (r GatePreparedRecord) Payload() gate.OpenPayload
Payload returns the private typed payload the resolver uses for validation/restore.
func (GatePreparedRecord) Prepared ¶
func (r GatePreparedRecord) Prepared() event.GatePrepared
Prepared returns the private prepared projection for the serializer to marshal.
type IdempotencyCollisionError ¶
type IdempotencyCollisionError struct {
ID string
Seq uint64 // the ledger sequence the ORIGINAL (colliding) record already occupies
}
IdempotencyCollisionError reports that a record's idempotency id already names a durable record in the log with a DIFFERENT persisted kind or payload — a genuine id collision (a bug or a forged retry), never a legitimate duplicate retry (which is always byte-identical to what is already durable). The append fails closed rather than silently accepting a differently-shaped record under a reused id.
func (*IdempotencyCollisionError) Error ¶
func (e *IdempotencyCollisionError) Error() string
type IdempotencyIndex ¶
type IdempotencyIndex struct {
// contains filtered or unexported fields
}
IdempotencyIndex tracks, for every idempotency id already durable in a session's log, the ledger sequence it occupies and the Fingerprint of what was persisted under it. A backend hydrates one from its full durable ledger (see the sessionstore package) before accepting new appends, then consults and updates it — via Check and Observe — under the same lock that already serializes its writes. It is NOT safe for concurrent use on its own; the caller's append-serializing lock is its only synchronization.
func NewIdempotencyIndex ¶
func NewIdempotencyIndex() *IdempotencyIndex
NewIdempotencyIndex returns an empty index ready for hydration.
func (*IdempotencyIndex) Check ¶
func (idx *IdempotencyIndex) Check(id string, fp Fingerprint) (seq uint64, duplicate bool, err error)
Check consults the index for id against the CANDIDATE fingerprint fp of a record about to be appended:
- id has never been observed: (0, false, nil) — the caller should proceed to durably append the record as new.
- id was observed with an IDENTICAL fingerprint: (seq, true, nil) — the caller should report AppendResult{Sequence: seq, Appended: false} WITHOUT appending.
- id was observed with a DIFFERENT fingerprint: (0, false, *IdempotencyCollisionError) — the caller must fail the append closed.
func (*IdempotencyIndex) Observe ¶
func (idx *IdempotencyIndex) Observe(id string, seq uint64, fp Fingerprint)
Observe records that id occupies seq with fingerprint fp, overwriting any prior entry for id. A backend calls it once per record while hydrating from history, and once more immediately after each new durable append lands.
type IdempotentJournal ¶
type IdempotentJournal interface {
SessionJournal
// AppendIdempotent behaves exactly like Append — same fencing, same errors —
// except a record whose IdempotencyID() already names a durable record with an
// IDENTICAL persisted kind+payload is detected and reported via
// AppendResult.Appended=false (carrying the ORIGINAL sequence) rather than
// durably appended a second time. A record whose id names a durable record with
// a DIFFERENT persisted kind or payload fails closed with a typed
// *IdempotencyCollisionError.
AppendIdempotent(ctx context.Context, rec JournalRecord) (AppendResult, error)
}
IdempotentJournal is the OPTIONAL extension a SessionJournal implementation may satisfy to deduplicate a redelivered append by idempotency id. It embeds SessionJournal so an idempotent implementation is usable anywhere a plain SessionJournal is expected — the existing narrow Append seam is never weakened or replaced. A caller that additionally wants to know whether ITS OWN call produced a new durable frame or deduplicated a retry (e.g. to skip a live broadcast for a duplicate) type-asserts for IdempotentJournal and calls AppendIdempotent instead of Append.
type JournalCommandAppender ¶
type JournalCommandAppender struct {
// contains filtered or unexported fields
}
JournalCommandAppender adapts a SessionJournal to the narrow "append one command" seam the session depends on for the intent log. The session holds an unexported commandAppender interface (AppendCommand(ctx, CommandRecord) error); this type satisfies it structurally, so the composition root (Phase 10) wires it in without the session importing journal internals beyond the CommandRecord constructor (Dependency Inversion). It carries no state beyond the journal — one method, one responsibility: append a CommandRecord (which the SESSION built with the dispatch target loopID, since a command — Interrupt/Shutdown especially — does not carry its own routing) and return the underlying typed error.
Unlike the event appender, the SESSION treats this seam as AUDIT-ONLY: a non-nil error is logged and the dispatch proceeds (losing a command record must never block the user's action). This façade itself never swallows — it returns the journal's error unchanged so the session owns the log-and-proceed decision.
func NewJournalCommandAppender ¶
func NewJournalCommandAppender(journal SessionJournal) *JournalCommandAppender
NewJournalCommandAppender wraps journal as a command appender. Like the event appender's unchecked form, it does NOT guard against a nil journal — use NewJournalCommandAppenderChecked at the composition root where a wiring bug must fail loud.
func NewJournalCommandAppenderChecked ¶
func NewJournalCommandAppenderChecked(journal SessionJournal) (*JournalCommandAppender, error)
NewJournalCommandAppenderChecked is the fail-loud constructor for the composition root: it returns a typed *NilJournalError if journal is nil rather than deferring the failure to a nil-deref at the first append.
func (*JournalCommandAppender) AppendCommand ¶
func (a *JournalCommandAppender) AppendCommand(ctx context.Context, rec CommandRecord) error
AppendCommand appends one intent-log command record: it calls the journal's Append with the session-built CommandRecord and returns the underlying typed error unchanged (the session logs+proceeds — audit-only — never faulting the session on a command-append failure). The CommandRecord routes to the target loop's command (intent-log) subject and uses the command's CommandID as the Nats-Msg-Id (idempotency). The returned sequence is discarded — the session needs only the success/failure signal.
type JournalEventAppender ¶
type JournalEventAppender struct {
// contains filtered or unexported fields
}
JournalEventAppender adapts a SessionJournal (the write side) to the narrow "append one Enduring event" seam the session hub depends on. The hub holds an unexported eventAppender interface (AppendEvent(ctx, event.Event) error); this type satisfies it structurally, so the composition root (Phase 10) wires it in via hub.WithAppender without the hub ever importing the journal package (Dependency Inversion). Beyond the journal it holds an optional catalog updater (nop by default): after a successful Append it notifies the catalog best-effort so the replay-free session index stays current. One responsibility: wrap the event in an EventRecord (which self-derives its subject from the event's scope+coordinates and its idempotency id from the EventID), append it, then best-effort index it.
func NewJournalEventAppender ¶
func NewJournalEventAppender(journal SessionJournal, opts ...AppenderOption) *JournalEventAppender
NewJournalEventAppender wraps journal as an event appender. It does NOT guard against a nil journal — use NewJournalEventAppenderChecked at the composition root where a wiring bug must fail loud. This unchecked form exists for call sites that have already validated the journal (and for the structural-satisfaction assertion).
func NewJournalEventAppenderChecked ¶
func NewJournalEventAppenderChecked(journal SessionJournal, opts ...AppenderOption) (*JournalEventAppender, error)
NewJournalEventAppenderChecked is the fail-loud constructor for the composition root: it returns a typed *NilJournalError if journal is nil rather than deferring the failure to a nil-deref at the first append.
func (*JournalEventAppender) AppendEvent ¶
AppendEvent durably appends one Enduring event: it wraps ev in an EventRecord and calls the journal's Append, returning the underlying typed error unchanged (the hub maps it onto a SessionPersistenceFault — never swallowed). The EventRecord routes a session-scoped event to the session subject and a loop-scoped event to its loop event subject, and uses the event's EventID as the Nats-Msg-Id (idempotency). An Ephemeral event is never appended by the hub; if one were passed, the serializer's event.MarshalEvent fails closed inside Append, so this path stays fail-secure.
It returns the assigned durable journal sequence so the hub can ride it on the LIVE delivery (event.Delivery.JournalSeq) — the sequence NEVER enters the persisted event codec. ONLY after the durable append succeeds does it best-effort notify the catalog (with the same sequence) so the replay-free session index stays current. The catalog update is the soft tail: its error is swallowed inside UpdateOnEvent and cannot change this method's return — the durable append stays strict, the catalog is derivable. On an append failure the catalog is NOT touched (the event did not durably land) and seq 0 is returned alongside the error.
When the underlying journal additionally satisfies IdempotentJournal (the optional dedup seam), a redelivered event whose EventID already names a durable record is detected there and reported as AppendResult.Appended=false; this method then returns the ORIGINAL sequence without re-notifying the catalog — a duplicate was already indexed by its first, genuine append, so republishing it a second time would be redundant. A journal that does not implement the optional interface behaves exactly as before (every successful Append notifies the catalog).
func (*JournalEventAppender) AppendEventResult ¶
func (a *JournalEventAppender) AppendEventResult(ctx context.Context, ev event.Event) (uint64, bool, error)
AppendEventResult is the result-preserving event append seam used by the Hub trusted publication path. Appended is true only when this call created a new durable frame; an identical idempotent retry returns the original sequence and Appended=false. The legacy AppendEvent method above deliberately discards only this boolean so existing callers retain their API and error behavior.
type JournalGateAppender ¶
type JournalGateAppender struct {
// contains filtered or unexported fields
}
JournalGateAppender adapts a SessionJournal to the session gate directory's strict durable append seam. GatePreparedRecord is appended as a private record; GateOpened and GateResolved are public Enduring events and are wrapped in EventRecord exactly like JournalEventAppender.
func NewJournalGateAppender ¶
func NewJournalGateAppender(journal SessionJournal) *JournalGateAppender
NewJournalGateAppender wraps journal as a gate appender. Like the other unchecked appender constructors, it expects a validated journal; use NewJournalGateAppenderChecked at composition roots.
func NewJournalGateAppenderChecked ¶
func NewJournalGateAppenderChecked(journal SessionJournal) (*JournalGateAppender, error)
NewJournalGateAppenderChecked fails loud on nil journal so composition wiring bugs surface at construction instead of the first gate operation.
func (*JournalGateAppender) AppendGateOpened ¶
func (a *JournalGateAppender) AppendGateOpened(ctx context.Context, ev event.GateOpened) error
func (*JournalGateAppender) AppendGatePrepared ¶
func (a *JournalGateAppender) AppendGatePrepared(ctx context.Context, rec GatePreparedRecord) error
func (*JournalGateAppender) AppendGateResolved ¶
func (a *JournalGateAppender) AppendGateResolved(ctx context.Context, ev event.GateResolved) error
type JournalLeaseLostError ¶
JournalLeaseLostError reports an Append refused because the journal's ownership lease was lost — released by the holder or overtaken by a higher epoch. Once the lease is gone the journal fails every append fast and never re-fetches or advances its expected sequence: a new owner (higher epoch) has written, or will write, its own LeaseFence, so this stale journal's fence would reject the append anyway. Failing here is the fast-path guard; the backend fence is the hard backstop. It carries the session and the lost lease's epoch and unwraps to a *LeaseLostError for errors.As.
func (*JournalLeaseLostError) Error ¶
func (e *JournalLeaseLostError) Error() string
func (*JournalLeaseLostError) Unwrap ¶
func (e *JournalLeaseLostError) Unwrap() error
type JournalNotReadyError ¶
JournalNotReadyError reports an Append attempted before the journal's opening LeaseFence was acknowledged. The journal writes the LeaseFence as its first append and only marks itself ready once it lands; an Append before that fails closed with this typed error rather than racing the fence. It carries the session so a caller can correlate the failure.
func (*JournalNotReadyError) Error ¶
func (e *JournalNotReadyError) Error() string
type JournalRecord ¶
type JournalRecord interface {
// IdempotencyID is the stable per-record id a backend uses as its de-dup key so
// a redelivered append de-duplicates: an event's EventID, a command's physical
// CommandRecordID, or a fence's epoch.
IdempotencyID() string
// contains filtered or unexported methods
}
JournalRecord is the sealed sum a session's serialized writer persists: an Enduring event, a command (the intent log), or an internal LeaseFence. It is a marker plus the one backend-neutral fact the writer needs to persist without re-inspecting the payload — the record's idempotency id (the stable per-record id a backend uses to de-duplicate a redelivered append). The concrete payload codec is the existing event/command marshaler; a record only carries the typed payload and exposes how to identify it. How a record is routed/stored is the backend's concern (a storage ledger name, a subject, …), never the record's.
The set is sealed by the unexported isJournalRecord marker: only the wrapper types in this package implement it, so the serializer's switch over the sum is exhaustive and a foreign type can never masquerade as a record.
type Lease ¶
type Lease interface {
// SessionID is the session this lease grants single-writer ownership of.
SessionID() uuid.UUID
// Release relinquishes the lease: stops any heartbeat, marks it no longer held
// (firing Lost), and best-effort clears the entry so a successor can re-acquire
// without waiting out the TTL. Idempotent.
Release(ctx context.Context) error
// contains filtered or unexported methods
}
Lease is the single-writer ownership token for one session's durable log. A SessionJournal depends on it (DIP): the composition root acquires a Lease from a backend and passes it in; the journal stamps the lease's Epoch into its first LeaseFence and refuses to append once the lease is lost. The holder (the composition root), not the journal, calls Release — the journal only reads Epoch and the validity/loss signals (see ownershipToken, the narrower view it depends on).
type LeaseFence ¶
type LeaseFence struct {
Epoch uint64 `json:"epoch"`
}
LeaseFence is an internal journal record marking a lease-handover boundary: a monotonically increasing Epoch fenced into the stream when ownership of a session's writer lease changes. It is journal-private (not an event or a command); the EventReplayer never decodes it. Codec: MarshalLeaseFence/UnmarshalLeaseFence in record_json.go.
func UnmarshalLeaseFence ¶
func UnmarshalLeaseFence(data []byte) (LeaseFence, error)
UnmarshalLeaseFence decodes bytes produced by MarshalLeaseFence. It fails closed with a *FenceDecodeError on any malformed input — empty bytes, non-object JSON, a non-numeric/negative epoch, an unknown field, or trailing bytes after the object. A negative or non-integer epoch fails because Epoch is a uint64.
type LeaseHeldError ¶
LeaseHeldError reports that acquiring a lease lost the single-holder race: the session's lease is currently held by a live (unexpired) holder, or a concurrent acquirer won the race. It carries the session and the epoch currently fenced so a caller can log who holds it. It is the expected, non-fatal "someone else owns this session" outcome — the loser must not write to the log.
func (*LeaseHeldError) Error ¶
func (e *LeaseHeldError) Error() string
type LeaseLostError ¶
LeaseLostError reports an operation attempted on a lease that is no longer held: it was released, or a higher-epoch holder took over (detected on a heartbeat renewal). It carries the session and the lease's epoch. The journal returns it (wrapped in a JournalLeaseLostError) when an Append is attempted after the lease is lost.
func (*LeaseLostError) Error ¶
func (e *LeaseLostError) Error() string
type MarshalRecordError ¶
MarshalRecordError wraps a failure to encode a record's payload before it is persisted. It names the record's routing/destination identifier so a caller can correlate the failure without re-inspecting the payload, and unwraps to the underlying codec error (an *event.EphemeralNotPersistableError, *command.UnknownCommandTypeError, a *FenceEncodeError, …) for errors.As inspection.
func (*MarshalRecordError) Error ¶
func (e *MarshalRecordError) Error() string
func (*MarshalRecordError) Unwrap ¶
func (e *MarshalRecordError) Unwrap() error
type NilJournalError ¶
type NilJournalError struct{}
NilJournalError reports that a JournalEventAppender was constructed over a nil SessionJournal — a composition-wiring bug. The checked constructor fails loud with this typed error rather than letting the nil surface as a panic at the first append.
func (*NilJournalError) Error ¶
func (*NilJournalError) Error() string
type RecordCursor ¶
type RecordCursor interface {
// Next returns the next record and its sequence, or io.EOF when the cold backlog is
// exhausted. A decode/read error fails secure: the cursor surfaces the typed error
// rather than skipping or zero-valuing the record.
Next(ctx context.Context) (JournalRecord, uint64, error)
// Close tears down the reader. Idempotent: a second call is a no-op.
Close() error
}
RecordCursor yields a session's journal records in sequence order. Next returns the next decoded JournalRecord (an EventRecord, CommandRecord, or FenceRecord) with its sequence, io.EOF once the cold backlog is drained, or a typed error on a malformed/missing/corrupt record. Close releases the underlying reader; it is idempotent and safe to call after an error. It is the all-records counterpart to EventCursor (which yields events only).
type RecordKindError ¶
type RecordKindError struct {
Subject string
}
RecordKindError reports a JournalRecord whose concrete type is outside the sealed sum the serializer encodes. It is unreachable for an in-package record (the sum is sealed by the unexported marker); it exists so a serializer's default arm fails closed with a typed error rather than panicking.
func (*RecordKindError) Error ¶
func (e *RecordKindError) Error() string
type RecordReplayer ¶
type RecordReplayer interface {
// Open binds a cursor over the WHOLE session log (every record kind) positioned at
// req.From. Only the cold path (Follow:false) need be implemented; Follow:true
// returns a typed *FollowUnsupportedError, matching EventReplayer.
Open(ctx context.Context, req ReplayRequest) (RecordCursor, error)
}
RecordReplayer is the journal's FULL read side: it opens an ordered cursor over a session's log and surfaces EVERY record — events, commands, AND fences — in sequence order. It is the data seam the transcript export consumes: the narrower EventReplayer yields enduring events only and therefore DROPS every CommandRecord (the user's gate decisions). Reading the whole log in sequence instead yields events and commands interleaved in append/causal order — the merged stream the transcript builder needs. The concrete implementation lives in a backend package (same inputs as EventReplayer).
type RecordTooLargeError ¶
RecordTooLargeError reports a record whose marshaled payload exceeded the inline threshold but could NOT be offloaded to the backend's content-addressed blob store (the store was unavailable or the upload failed). The journal fails closed with this typed error rather than silently persisting an over-threshold record. It carries the record's routing/destination identifier, its idempotency id, the payload length, and the underlying offload cause.
func (*RecordTooLargeError) Error ¶
func (e *RecordTooLargeError) Error() string
func (*RecordTooLargeError) Unwrap ¶
func (e *RecordTooLargeError) Unwrap() error
type ReplayRequest ¶
type ReplayRequest struct {
// SessionID is the session whose log is replayed (required; a zero id yields a
// setup error rather than a replay over every session).
SessionID uuid.UUID
// LoopID, when non-zero, narrows an event replay to that single loop; zero
// replays the session's events plus every loop's events.
LoopID uuid.UUID
// From is where the backlog read begins: Beginning or FromSeq(n).
From StartPos
// Follow keeps the cursor live after the backlog drains (tailing new appends). A
// backend that implements only the cold path returns a typed *FollowUnsupportedError
// from Open rather than silently behaving as a cold cursor.
Follow bool
}
ReplayRequest selects which of a session's records to replay and how. Which records (events only, or a single loop's) is derived from SessionID + LoopID; how far back from From; whether to keep tailing from Follow. The concrete filtering is the backend replayer's job — this is the backend-neutral request it honors.
type SessionJournal ¶
type SessionJournal interface {
// Append serializes rec, persists it under the next expected sequence, and
// returns the assigned sequence. ctx bounds the caller's willingness to wait; the
// implementation additionally carries a per-append deadline independent of ctx so
// one stuck call cannot wedge the serialized writer forever. Appends are totally
// ordered: the returned sequences are strictly monotonic across calls.
Append(ctx context.Context, rec JournalRecord) (seq uint64, err error)
}
SessionJournal is the single serialized writer for one session's durable log. Append encodes a JournalRecord's payload, persists it under single-writer fencing, and returns the assigned sequence. It is the only thing that writes a session's log; callers funnel every event, command, and fence through it so the log stays a totally-ordered, gap-free record of the session.
The interface is intentionally narrow (one method): a caller that only needs to persist a record must not depend on any log-management surface. The concrete implementation lives in a backend package (e.g. pkg/sessionstore over storage), wired at the composition root — this package owns only the contract.
func WithHooks ¶
func WithHooks(j SessionJournal, runner *hook.Runner, sessionID uuid.UUID) SessionJournal
WithHooks observes each durable append while preserving the journal's result.
type StartPos ¶
type StartPos struct {
// contains filtered or unexported fields
}
StartPos is the closed value type naming where a replay begins: the log beginning (every record) or a specific sequence. It is a value, not an interface, so a caller cannot smuggle a third start mode past a switch; the two constructors Beginning and FromSeq are the only ways to build one, and Seq reads it back.