Documentation
¶
Overview ¶
Package server implements the Bolt v5 TCP server for the GoGraph Cypher engine. It handles connection acceptance, Bolt protocol negotiation, session lifecycle, and authentication.
Concurrency ¶
Server is safe for concurrent use by multiple goroutines. Session and State are NOT safe for concurrent use; each connection owns exactly one Session.
Index ¶
- Constants
- Variables
- func ConstantTimeValidate(wantPrincipal, wantCredentials string) func(principal, credentials string) error
- func DefaultTLSConfig() *tls.Config
- func ExtractBookmarks(extra map[string]packstream.Value) []string
- func FailureCode(err error) string
- func NextBookmark() string
- func RoutingTable(addr string) map[string]packstream.Value
- type AuthHandler
- type BasicAuthHandler
- type CertReloader
- type Identity
- type NegativeTimeoutError
- type NoAuthHandler
- type Options
- type Server
- func (s *Server) ListenAndServe(ctx context.Context, addr string) error
- func (s *Server) Serve(ctx context.Context, ln net.Listener) (err error)
- func (s *Server) SetClock(clk clock.Clock)
- func (s *Server) Shutdown(ctx context.Context) error
- func (s *Server) TerminateTransaction(id string) error
- func (s *Server) Transactions() []TransactionInfo
- type Session
- type State
- type TransactionInfo
- type Tx
Examples ¶
Constants ¶
const ( // DefaultMaxInFlightPerConnection is the default value applied to // Options.MaxInFlightPerConnection when the caller leaves it at // zero. The count tracks all Result cursors appended to the // in-progress explicit transaction since BEGIN (both open and // already-drained), so it bounds the total number of RUN statements // a client may issue without committing. The default of 1024 allows // any legitimate workload while still bounding pathological // RUN-loop attacks that grow tx.results without bound. Operators // that need a stricter limit may lower this value explicitly. DefaultMaxInFlightPerConnection = 1024 // DefaultConnTimeout is the value applied to Options.ConnTimeout when the // caller leaves it at zero. It IS zero: the per-message idle read deadline is // DISABLED by default. See [Options.ConnTimeout] for what the field does when // an operator sets it, and for the three-way contract (0 disables, a positive // value bounds, a negative value is rejected by [NewServer]). // // # Why it is disabled (rmp #2807) // // The default was 30 s, and the bound was a read deadline armed before EVERY // read of the post-handshake message loop. The reader goroutine sits in that // read while the message loop executes the client's own statement, so the // deadline ran against a BUSY server rather than against an idle client: a // client waiting for the records it just asked for sends nothing, and this // deadline could not tell that apart from silence. Measured under rmp #2806: // with ConnTimeout 100 ms and a one-hour statement bound, an autocommit // statement of about 900 ms had its connection closed at 110 ms — the server // cutting itself off for being busy. Neither peer this module is measured // against does that. // // Both peers ship the equivalent bounds disabled. PostgreSQL ships // statement_timeout, transaction_timeout, idle_in_transaction_session_timeout // and idle_session_timeout all at 0 and advises against setting them globally // in postgresql.conf; its liveness check is TCP keepalive, never a read // deadline. Neo4j ships db.transaction.timeout at 0 and keeps connections // live with protocol-level NOOP chunks. // // # What replaces it // // TCP keepalive, enabled on every accepted connection — see // [DefaultKeepAliveIdle] for the budget and for what it does and does not // detect. A peer that is GONE (a pulled cable, a NAT that dropped its state, // a host that vanished without sending FIN) is reclaimed by the kernel within // that budget, which is what the read deadline was being misused to // approximate. // // The bounds on COUNT are untouched and are what bound the resource: // [Options.MaxConnections] (1024), [DefaultMaxOpenTxPerPrincipal] (2048) and // the reclamation horizon's fixed slot count // (graph/mvcc.HorizonCapacity, 1024). PostgreSQL bounds by count too; only // the time bound goes. // // THE COST, stated rather than glossed, in two parts: // // - A client that is ALIVE and silent — its TCP stack answers keepalive // probes, its process simply never sends another Bolt message — is no // longer reclaimed at all. It holds one of the 1024 connection slots and // its two goroutines until it disconnects or the server shuts down, and // where it left an explicit transaction open, that transaction pins the // MVCC reclamation horizon and one of its 1024 slots for exactly as long. // The remedy is the operator's, and it is the one PostgreSQL offers for // the same exposure: set [Options.MaxTxIdleTime] to reclaim the // transaction and [Options.ConnTimeout] to reclaim the connection. // - The window between a completed handshake and a successful LOGON is no // longer bounded either. [DefaultHandshakeTimeout] bounds version // negotiation only, so an unauthenticated client that negotiates and then // falls silent now holds its slot indefinitely. PostgreSQL bounds that // window with a separate authentication_timeout (1 minute by default); // GoGraph has no such bound today. An operator exposing this server to an // untrusted network should set Options.ConnTimeout explicitly until it // does. DefaultConnTimeout time.Duration = 0 // DefaultTxTimeout is the value applied to Options.DefaultTxTimeout when the // caller leaves it at zero. It IS zero: an explicit transaction that supplies // no tx_timeout of its own gets NO total wall-clock bound by default. // // When an operator does set the field, it bounds an explicit transaction // (opened by BEGIN) however busy it is — [Options.MaxTxIdleTime] is the // separate bound on silence — and a client-supplied tx_timeout still takes // precedence, with Options.MaxStatementTimeout clamping the effective value // when it is set. // // # Why it is disabled (rmp #2807) // // The default was 30 s until rmp #2806 and 30 minutes after it. It is now 0, // matching PostgreSQL's transaction_timeout, which ships at 0 with the // documentation advising against setting it globally, and Neo4j's // db.transaction.timeout, which also ships at 0. // // A bounded default was defensible when an explicit transaction held the // engine's writer serialisation from BEGIN until COMMIT/ROLLBACK: one stalled // client then blocked every other writer on the server, a liveness denial of // service (#1302). rmp #2305 retired that hold. An unfinished transaction now // costs memory rather than other clients' progress: it pins the reclamation // horizon, so no version it could still read is freed while it lives, and it // occupies one of the horizon's fixed slots (graph/mvcc.HorizonCapacity). // // THE COST, stated rather than glossed: a long-running transaction is no // longer killed for being long. A client looping RUN/PULL forever, or holding // a transaction open across unbounded think-time, holds its versions and its // horizon slot for as long as its connection lives. Nothing else reclaims it: // the idle bound is disabled too, and so is the connection read deadline. An // operator who wants a total bound sets Options.DefaultTxTimeout explicitly. DefaultTxTimeout time.Duration = 0 // DefaultMaxTxIdleTime is the value applied to Options.MaxTxIdleTime when the // caller leaves it at zero. It IS zero: an OPEN explicit transaction that // stops receiving messages is NOT reclaimed by default. // // When an operator does set the field, it bounds how long an open explicit // transaction may go without the client sending a message, which is a // different bound from [Options.DefaultTxTimeout]: that one caps a // transaction's total life however busy, while this one reclaims one that has // been ABANDONED. Every inbound message pushes this deadline forward; a silent // client pushes nothing. // // # What it protected, and what changed // // Before rmp #2305 an open transaction held the global visibility barrier, so // one authenticated client that sent BEGIN and stopped talking turned a 4.7 ms // read on every other connection into a 30.001 s stall followed by a hard // TransactionTimedOut, repeatable indefinitely (rmp #2175). That outage is // GONE: an open transaction holds no barrier and no writer serialisation, and // the gates in this package's e2e_concurrent_write_tx_test.go assert exactly // that against the official driver. // // What an abandoned transaction still costs is memory and horizon slots. It // pins the reclamation horizon, so no version it could still read is freed // while it lives, and it occupies one of the horizon's fixed slots // (graph/mvcc.HorizonCapacity). An availability failure became a // memory-and-slot failure. // // # Why it is disabled (rmp #2807) // // The default was 5 s until rmp #2806 and 30 minutes after it. It is now 0, // matching PostgreSQL's idle_in_transaction_session_timeout — the bound that // governs exactly this case, shipped at 0, with the documentation advising // against setting it globally in postgresql.conf. // // THE COST, stated rather than glossed: ONE abandoned transaction holds every // version it could still read, and one of the horizon's 1024 slots, for as // long as its connection lives — no longer for 5 s, nor for half an hour, but // indefinitely. N abandoned transactions hold N slots on the same terms, so // the horizon's capacity is exhaustible by silent clients alone, and // [Options.MaxConnections] (1024) is what bounds how many of them there can // be. Nothing reclaims them on a timer any more: [DefaultTxTimeout] is // disabled, and so is [DefaultConnTimeout], which used to tear the connection // down and roll the transaction back as a side effect. // // The remedy is the operator's, and it is PostgreSQL's posture exactly: set // Options.MaxTxIdleTime explicitly. Server.Transactions and // Server.TerminateTransaction remain the manual route for one transaction at // a time. // // The concurrency gates in e2e_concurrent_write_tx_test.go set ten minutes of // their own. That override is kept because it pins the value those gates // measure under instead of inheriting whatever this constant becomes. DefaultMaxTxIdleTime time.Duration = 0 // DefaultMaxOpenTxPerPrincipal is the value applied to // Options.MaxOpenTxPerPrincipal when the caller leaves it at zero. It caps // how many explicit transactions one authenticated principal may hold open at // once, across all of its connections. // // # It binds on WRITE transactions too as of rmp #2305 // // This note used to say the bound could never be the binding constraint for // write transactions, because one held the engine's writer serialisation for its // whole life and the engine therefore capped concurrently-open write // transactions at ONE server-wide. rmp #2305 retired that hold. Write // transactions now overlap freely, so this is a REAL limit for them, and a // client that opens more than the default 16 at once is refused with // LimitExceeded. Two goroutine-leak tests in this package had to raise it for // exactly that reason. // // A READ transaction (BEGIN with mode "r") has always been concurrent and // unbounded per principal without this; that is what it originally capped, along // with the session and cursor state each one holds — and, since rmp #2307, the // MVCC read snapshot each one pins for its lifetime, which holds the reclamation // horizon back until the handle finishes. A write transaction pins the horizon // the same way since rmp #2305, so both modes now cost the same thing. // // The count is of OPEN transactions, not of BEGINs waiting to be admitted: a // burst of concurrent BEGINs from one principal is bounded by MaxConnections, // not by this. Counting waiting BEGINs here would reject legitimate concurrent // traffic. // // # Why 2048, and what it gives up (rmp #2419) // // It was 16, defended here on the grounds that "a pool with more than sixteen // connections per principal should say so explicitly". The 2026-08-11 // concurrency assessment recorded the consequence as finding F2: CLAUDE.md // publishes 1, 8, 64, 256 and 1024 goroutines as the levels this module // measures and reports at, and a single principal could not reach them through // explicit transactions without overriding this first. Every harness that got // there had already had to — bench/soak with 1200, bench/comparison/ggserver // with a flag — which is the shape of a default disagreeing with a published // contract rather than of two harnesses being unusual. // // 2048 is above the highest published level, so the default configuration now // reaches the concurrency the module publishes. Note what that means in // practice: a connection holds at most one open transaction and // [Options.MaxConnections] defaults to 1024, so under the default // configuration this quota CANNOT BIND — the connection ceiling is reached // first, and it is the connection ceiling that bounds the resource. This // quota binds again only for an operator who raises MaxConnections above // 2048, and it remains what isolates one principal from another. // // THE COST, stated rather than glossed: every open transaction pins an MVCC // read snapshot and holds the reclamation horizon back for its lifetime // (rmp #2305, #2307), so a higher ceiling is a weaker bound on that resource. // NOTHING limits the damage on a timer any more. [DefaultMaxTxIdleTime] was // 5 s when this ceiling was raised and 30 minutes after rmp #2806; since // rmp #2807 it is 0, so an abandoned transaction is held until its client // disconnects. This ceiling and [Options.MaxConnections] are therefore the // bounds that actually apply, and an embedder that wants a time bound as well // sets Options.MaxTxIdleTime explicitly. One that wants the old tight count // bound sets Options.MaxOpenTxPerPrincipal explicitly, which is the only way // to get it. // // Set a negative value to disable enforcement — which is a deliberate, // visible choice at the call site, not something reachable by accident. DefaultMaxOpenTxPerPrincipal = 2048 // DefaultDatabaseName is the value applied to Options.DatabaseName when the // caller leaves it empty, and the name reported in the `db` field of result // metadata for a client that selected no database. // // It is "neo4j" because that is the database name every Bolt client assumes // when it is given none — the driver's own default, and what cypher-shell and // the Neo4j Browser display. Reporting it is a compatibility choice about a // label, not a claim about the product: the server identifies itself as // GoGraph in the `server` and `bolt_agent` fields of the HELLO response. // Operators serving a differently named graph may override it. // // GoGraph serves one graph per server, so the name selects nothing. An // unknown name from a client is echoed rather than rejected: Neo4j answers a // missing database with Neo.ClientError.Database.DatabaseNotFound, which // would be the stricter behaviour, but rejecting a name the server has always // accepted would break existing embedders and is a separate decision from // reporting the field at all. DefaultDatabaseName = "neo4j" // DefaultStatementTimeout is the value applied to // Options.DefaultStatementTimeout when the caller leaves it at zero. It IS // zero: an AUTOCOMMIT statement (a bare RUN outside an explicit transaction) // that supplies no per-statement timeout of its own gets NO wall-clock bound // by default. // // When an operator does set the field, a client-supplied `timeout` still takes // precedence, and Options.MaxStatementTimeout, when set, additionally clamps // the effective value. // // # Why it is disabled (rmp #2807) // // The default was 30 s until rmp #2806 and 30 minutes after it. It is now 0, // matching PostgreSQL's statement_timeout, which ships at 0 with the // documentation advising against setting it globally in postgresql.conf. // // THE COST, stated rather than glossed: a default-configured server has no // wall-clock bound on an autocommit statement, which is the state #1828 // recorded as a defect when the alternative was an unbounded bound on an // otherwise bounded server. An authenticated client can submit a statement // whose runtime is super-linear in the graph size yet whose result collapses // to a single row — a disconnected multi-pattern Cartesian product such as // `MATCH (a),(b),(c),(d),(e) RETURN count(*)` — so the result-row and // result-byte caps never fire and the statement pins a CPU core until it // finishes or the client disconnects. [DefaultConnTimeout] no longer cuts it // short either, because it too is disabled. // // What still bounds it: the engine's own result caps // (cypher.EngineOptions.MaxResultRows and MaxResultBytes) for any statement // whose result actually grows, the connection ceiling for how many such // statements can run at once, and cancellation — a client that disconnects // cancels the connection context, which stops the statement. An operator who // wants a wall-clock bound sets Options.DefaultStatementTimeout, or the // server-wide Options.MaxStatementTimeout, explicitly. DefaultStatementTimeout time.Duration = 0 // DefaultHandshakeTimeout is the deadline that bounds the unauthenticated // version-negotiation handshake — the cheapest phase for an attacker to // abuse, since it requires no valid protocol bytes (a client may open a // socket, send a single byte, and otherwise stall). The deadline is applied // to the connection before [proto.Negotiate] and cleared on success so it // never bleeds into normal operation. A legitimate client sends its 20-byte // handshake immediately, so 10 s is ample, while a stalled handshake is // reclaimed promptly. The handshake bound is fixed (not configurable via // Options) to keep the Options struct small; the package var handshakeTimeout // is seeded from this const and overridable only by tests. // // It bounds VERSION NEGOTIATION ONLY. Since rmp #2807 disabled // [DefaultConnTimeout], nothing bounds the window between a completed // handshake and a successful LOGON by default, so this is no longer the // shorter of two pre-authentication bounds — it is the only one. See // [DefaultConnTimeout] for what that costs and how an operator closes it. DefaultHandshakeTimeout = 10 * time.Second // DefaultKeepAliveIdle is how long an accepted connection must be silent // before the kernel sends its first TCP keep-alive probe. // // It is one of three parameters — with [DefaultKeepAliveInterval] and // [DefaultKeepAliveCount] — that the server installs on every accepted // connection that is a *[net.TCPConn]. Together they bound how long a // connection whose peer has GONE (a pulled cable, a NAT or firewall that // silently dropped its state, a host that vanished without sending FIN) // survives: the first probe after this idle period, further probes every // interval, and the connection dropped after that many consecutive probes go // unanswered. 15 s + 3 x 5 s is 30 s, which is deliberately the reclamation // budget the 30 s [DefaultConnTimeout] used to give a dead connection before // rmp #2807 disabled it. // // This is the liveness mechanism PostgreSQL relies on, and it REPLACES the // read deadline rather than supplementing it. What it detects is a DEAD peer. // What it does NOT detect is a live client that is merely silent: its stack // answers every probe, so nothing here reclaims it. That is the exposure // [DefaultConnTimeout] documents, and Options.ConnTimeout is what closes it. // // The values are fixed rather than configurable, for the same reason // [DefaultHandshakeTimeout] is: the Options struct stays small. Go's own // listener default (15 s / 15 s / 9, so roughly 150 s) applies to a // connection accepted from a net.Listener the caller did not configure; the // server overrides it on every accepted connection so the budget does not // depend on how the embedder built the listener, nor on the platform default // (7200 s / 75 s / 8, measured on this project's development host) when the // embedder disabled the listener's own keep-alive. DefaultKeepAliveIdle = 15 * time.Second // DefaultKeepAliveInterval is the gap between successive TCP keep-alive // probes on an accepted connection. See [DefaultKeepAliveIdle]. DefaultKeepAliveInterval = 5 * time.Second // DefaultKeepAliveCount is how many consecutive TCP keep-alive probes may go // unanswered before the kernel drops the connection. See // [DefaultKeepAliveIdle]. DefaultKeepAliveCount = 3 )
const DefaultMaxInboundDecodeBytes int64 = 1 << 30 // 1 GiB
DefaultMaxInboundDecodeBytes is the engine-wide inbound-decode ceiling applied when Options.MaxInboundDecodeBytes is left at zero AND the process has no Go soft memory limit to derive one from.
It exists because the GOMEMLIMIT derivation below silently produced NO ceiling in the commonest deployment. An unset GOMEMLIMIT is the Go runtime's default — debug.SetMemoryLimit(-1) then reports math.MaxInt64 — so the documented "engine-wide inbound-memory ceiling" was inert unless the operator had separately set a memory limit, and the real bound was MaxConnections times the per-connection limits: 1024 x 16 MiB of reassembly buffers plus 1024 x 128 MiB of decoded collections. That allocation is reachable PRE-AUTHENTICATION, because a HELLO must be decoded before it can be authenticated. The bounded-resources mandate requires an explicit finite upper bound, and the existence of MaxInboundDecodeBytesUnlimited settles that zero was never meant to mean unlimited: an opt-out sentinel would be redundant if it did.
The value is 1 GiB, matching github.com/FlavioCFOliveira/GoGraph/cypher.DefaultMaxResultBytes so the module's finite defaults share one scale. It is 64x the largest single message a client may send (proto.DefaultMaxMessageBytes, 16 MiB) and 8x the worst-case decoded size of one message, so it admits several concurrent large-message decodes while bounding a hostile fleet's aggregate to something any production host survives. A deployment that sets GOMEMLIMIT keeps the derived one-eighth fraction and is unaffected.
const MaxInboundDecodeBytesUnlimited int64 = -1
MaxInboundDecodeBytesUnlimited is the explicit opt-out sentinel for Options.MaxInboundDecodeBytes: set the field to this value to disable the engine-wide inbound-decode ceiling entirely. It is distinct from the zero value, which selects the GOMEMLIMIT-derived default.
Variables ¶
var ( // ErrAuthFailed is returned when credentials are invalid. ErrAuthFailed = errors.New("bolt: authentication failed") // ErrSchemeUnknown is returned when the auth scheme is not supported. ErrSchemeUnknown = errors.New("bolt: unknown auth scheme") )
Common auth errors.
var ErrCertOutsideValidity = errors.New("bolt/server: certificate outside its validity window")
ErrCertOutsideValidity is the sentinel wrapped by every CertReloader.Reload refusal caused by the validity window of the leaf on disk: it has expired, or it is not valid yet. Callers that distinguish a doomed renewal from unreadable or mismatched material — an alerting hook, say — match it with errors.Is.
It is NOT returned for a pair that fails to read, parse or pair; those keep their own wrapped errors.
var ErrInvalidTransition = errors.New("bolt: invalid state transition")
ErrInvalidTransition is returned by Transition when the given message type is not permitted in the current state.
var ErrNegativeTimeout = errors.New("bolt: negative timeout")
ErrNegativeTimeout is the sentinel every negative-duration refusal from NewServer unwraps to, so a caller can classify one with errors.Is(err, ErrNegativeTimeout) without naming each field. The concrete error is a *NegativeTimeoutError, which carries the field and the value; reach it with errors.As.
var ErrNoAuthHandler = errors.New("bolt: no auth handler configured; set Options.Auth to a real AuthHandler, or to NoAuthHandler{} to run without authentication")
ErrNoAuthHandler is returned by NewServer when Options.Auth is nil. The server is secure-by-default: running without authentication must be an explicit opt-in, never an accidental default. Set Options.Auth to a real AuthHandler to require credentials, or set Options.Auth to a NoAuthHandler{} value to run the open-door handler on purpose (development and testing only).
var ErrNoSuchTransaction = fmt.Errorf("bolt: no such open transaction")
ErrNoSuchTransaction is returned by Server.TerminateTransaction for an ID that is not open — either never seen, or already ended between a listing and the termination call.
Functions ¶
func ConstantTimeValidate ¶
func ConstantTimeValidate(wantPrincipal, wantCredentials string) func(principal, credentials string) error
ConstantTimeValidate returns a Validate function that accepts only the given principal and credentials, using crypto/subtle.ConstantTimeCompare for both comparisons. The comparison time is independent of the values being compared, eliminating timing side-channels.
Example:
handler := server.BasicAuthHandler{
Validate: server.ConstantTimeValidate("alice", "correct-horse-battery-staple"),
}
func DefaultTLSConfig ¶
DefaultTLSConfig returns a hardened baseline tls.Config that operators should use as the STARTING POINT for the server's transport security.
The configuration sets a TLS 1.2 floor and a modern, AEAD-only cipher list for the TLS 1.2 handshake; TLS 1.3 is negotiated automatically when both peers support it (TLS 1.3 cipher suites are fixed by the Go runtime and are always safe, so they are not — and cannot be — listed here). No MaxVersion is set, so a 1.3-capable client always upgrades to 1.3.
The returned config is INCOMPLETE on its own: it carries no certificate. Callers MUST populate it with their own server identity before use, by setting one of:
- tls.Config.Certificates, or
- tls.Config.GetCertificate (for example by wiring the on-disk hot reloader: cfg.GetCertificate = reloader.GetCertificate; see CertReloader).
Then pass the result in Options.TLSConfig.
The server does NOT impose this baseline automatically. It wraps whatever Options.TLSConfig the operator supplies verbatim: passing a nil TLSConfig keeps the existing behaviour of running PLAINTEXT TCP (no TLS at all). DefaultTLSConfig only provides and documents a safe default; it never overrides an operator-supplied config, so embedders are never surprised.
A fresh, independent config is returned on every call (no shared mutable global), so callers may freely mutate the result — adding Certificates, GetCertificate, client-auth policy, etc. — without aliasing another caller's configuration.
func ExtractBookmarks ¶
func ExtractBookmarks(extra map[string]packstream.Value) []string
ExtractBookmarks returns the bookmark list from RUN/BEGIN extra metadata. It reads the "bookmarks" key, which may be a []packstream.Value of strings. Returns nil (not an error) when the key is absent or the value is not a list.
func FailureCode ¶
FailureCode returns the Neo4j-style dot-delimited error code for err. Falls back to "Neo.DatabaseError.General.UnknownError" for unrecognised errors. The lookup uses errors.As and errors.Is so wrapped errors are matched correctly.
func NextBookmark ¶
func NextBookmark() string
NextBookmark generates a new bookmark string for a committed transaction. The format is "FB:kXXXXXX" where XXXXXX is a monotonically increasing counter expressed as a zero-padded 8-digit hexadecimal value.
NextBookmark is safe for concurrent use.
func RoutingTable ¶
func RoutingTable(addr string) map[string]packstream.Value
RoutingTable returns the single-host routing table for the server at addr. The TTL is hardcoded to 300 seconds. All three roles (WRITE, READ, ROUTE) point to the same single-host address.
The returned map matches the Bolt v5 routing table format expected inside a SUCCESS metadata "rt" key.
Types ¶
type AuthHandler ¶
type AuthHandler interface {
// Authenticate validates the auth scheme, principal, and credentials.
// On success it returns an Identity; on failure it returns a non-nil error.
// Returning ErrAuthFailed causes the server to send a Failure with code
// "Neo.ClientError.Security.Unauthorized". Returning ErrSchemeUnknown
// causes a Failure with code "Neo.ClientError.Security.AuthProviderFailed".
Authenticate(scheme, principal, credentials string) (Identity, error)
}
AuthHandler is the pluggable authentication interface. Implementations must be safe for concurrent use.
type BasicAuthHandler ¶
type BasicAuthHandler struct {
// Validate is called with the principal and credentials from the client.
// It must return nil on success and a non-nil error on failure.
// See the type-level documentation for timing side-channel guidance.
Validate func(principal, credentials string) error
}
BasicAuthHandler validates credentials by delegating to a caller-supplied Validate function. The Validate function must return nil on success and a non-nil error (typically ErrAuthFailed) on failure.
Timing side-channels ¶
Validate is called with the raw credential string from the client. If the implementation compares credentials with == or strings.Equal, an attacker can infer the correct value by measuring response latency (timing side-channel). Always use ConstantTimeValidate or crypto/subtle.ConstantTimeCompare for credential comparison:
handler := BasicAuthHandler{
Validate: ConstantTimeValidate("alice", "correct-horse-battery-staple"),
}
Do not add rate-limiting or account-lockout logic inside Validate; place it in a middleware wrapping the AuthHandler instead so the Bolt server remains stateless per-connection.
Attempt limits: what the server does and does not bound ¶
Because lockout lives outside this interface, the server itself bounds credential attempts only in one place. A FIRST authentication that fails terminates the connection (the session goes DEFUNCT), so a client gets one guess per connection before it must dial again. A RE-authentication — a LOGON on a connection that already authenticated — is NOT bounded: it fails the session, RESET recovers it, and the client may try again on the same connection indefinitely. Measured, not inferred: five wrong-password LOGONs on one connection, each recovered by RESET (rmp #2481).
That asymmetry is a consequence of the design above, not an oversight, and it is stated here so it is a conscious contract: an embedder who needs an attempt budget must impose it in the middleware wrapping this handler, where the per-principal state such a budget requires can live.
BasicAuthHandler is safe for concurrent use as long as Validate is.
func (BasicAuthHandler) Authenticate ¶
func (h BasicAuthHandler) Authenticate(scheme, principal, credentials string) (Identity, error)
Authenticate implements AuthHandler. It accepts only the "basic" scheme; any other scheme returns ErrSchemeUnknown. It calls h.Validate with the principal and credentials; if Validate returns a non-nil error, Authenticate returns ErrAuthFailed.
type CertReloader ¶
type CertReloader struct {
// contains filtered or unexported fields
}
CertReloader watches a (certificate, key) PEM file pair on disk and serves the most recent successfully loaded pair via the CertReloader.GetCertificate hook installable on tls.Config.GetCertificate.
The intent is operational: rotate the server's TLS material (e.g. cert-manager / Let's Encrypt) without restarting the Bolt server. The previous certificate stays in service until the new pair is fully validated and only then is the swap performed atomically via sync/atomic.Pointer. "Fully validated" means the pair reads, parses, pairs, AND its leaf is valid at the current instant: a reload that fails any of those leaves the live certificate untouched and surfaces the error via the provided OnError callback (or via stderr when nil). See CertReloader.Reload for the exact rule and ErrCertOutsideValidity.
CertReloader is safe for concurrent use; the hot path is a single atomic.Pointer.Load.
func NewCertReloader ¶
func NewCertReloader(certPath, keyPath string, onError func(error)) (*CertReloader, error)
NewCertReloader loads the certificate + key from disk and returns a CertReloader holding the result. The initial load is mandatory: if the files cannot be read or parsed, NewCertReloader returns the error and the caller MUST fail fast (do not start the server with a broken TLS config).
onError is invoked when a later reload (triggered by Reload or by the optional Watch goroutine) is REFUSED — whether because the new pair cannot be read, parsed or paired, or because its leaf is outside its validity window (see ErrCertOutsideValidity). A nil onError defaults to printing to stderr via fmt.Fprintln.
The initial load performed here is gated on read, parse and pair only; see CertReloader.Reload for why the validity check applies to the SWAP and not to construction.
func (*CertReloader) GetCertificate ¶
func (r *CertReloader) GetCertificate(_ *tls.ClientHelloInfo) (*tls.Certificate, error)
GetCertificate is the hook to install on tls.Config.GetCertificate. It returns the most recently loaded certificate. The signature matches the standard library's expectation so callers can do:
cfg := &tls.Config{GetCertificate: reloader.GetCertificate}
The returned *tls.Certificate is shared across all concurrent handshakes; callers must NOT mutate the returned value.
func (*CertReloader) Reload ¶
func (r *CertReloader) Reload() error
Reload re-reads the certificate + key from disk and atomically swaps the live certificate when the new pair is fully validated. A failure leaves the live certificate untouched and returns the error so the caller (or the OnError callback installed via NewCertReloader) can record the incident.
Validation is two-part. The pair must read, parse and pair — that is tls.X509KeyPair — and, when a certificate is already in service, the new leaf must also be valid AT THE CURRENT INSTANT. A leaf that has expired, or whose NotBefore is still in the future, is refused with an error wrapping ErrCertOutsideValidity and the live certificate stays in service. Without that second part a routine renewal that produced expired material — or a clock skew on the renewing host — would replace a working certificate with one that fails every handshake, converting a recoverable condition into a total outage of the listener (rmp #2557).
Why the skip is keyed on content, not on mtime ¶
Reload reads both files on every call and compares a SHA-256 digest of each against the digests the live certificate was built from; it re-parses only when they differ.
The skip is kept because it is measurably cheaper, though by less than intuition suggests: 20.6 µs and 2.0 kB per call against 37.8 µs and 9.4 kB for the full load, best of 5 on darwin/arm64 (BenchmarkCertReloader_ReloadUnchanged against BenchmarkCertReloader_ParseCostAvoidedBySkip). The two file reads dominate BOTH paths, so the skip is a 1.8x saving on a call that happens once per poll interval — it earns its place by not recomputing and re-publishing a certificate that has not changed, not by being fast. What matters far more is that it is PROVABLY a no-op.
The original skip compared file mtimes, and mtime is not a content hash. Any rotation that does not advance it — a rename from another directory, cp -p, a restore from an archive, or two rotations inside one filesystem timestamp tick — was reported as a SUCCESS having loaded nothing, so a rotation performed to REVOKE a certificate could be silently ignored while the operator believed the material had been replaced, and, because the call succeeded, onError never fired to say otherwise (rmp #2558).
The digests are recorded only on a load that succeeded, so a refused pair is always re-examined on the next call.
Three properties of that refusal matter operationally:
It does NOT record the refused pair's digests, so the very next Reload re-parses the same files. A pair refused for being not-yet-valid is therefore picked up by the Watch poller as soon as its NotBefore passes, with no operator action.
It checks the LEAF only, never the rest of the chain. An expired issuer deliberately kept in a bundle is a real, working deployment pattern — the Let's Encrypt DST Root CA X3 cross-sign after 2021-09-30 is the canonical case, and it was their DEFAULT recommended chain — and refusing it would break rotations that serve clients fine.
The residual is accepted deliberately, not overlooked: an expired LOAD-BEARING intermediate still breaks every handshake and this check will not catch it. Envoy Gateway is the one surveyed implementation that does validate the whole serving chain, and it documents two incidents its strictness caused (envoyproxy/gateway#9225 and #9473 — a configuration stall, and one broken Secret silently breaking healthy Gateways), because dropping an expired chain member leaves the key matching a certificate no longer served. Closing this residual by whole-chain validation would trade a narrow gap for a wider blast radius and for false refusals of chains that work.
It gates the SWAP, not the initial load. The contract being protected is that a WORKING certificate is never replaced by a broken one; at construction there is nothing to protect, and refusing there would make a host whose clock has not yet synchronised unable to start at all, which is strictly worse than serving material its clients may well accept. That is not a hypothetical: refusing a certificate for its dates AT START is a repeatedly reproduced outage (docker/for-win#2913, kubernetes/minikube#13779, k3s-io/k3s#6152, and the 2012 Azure leap-day disruption, where an agent that failed to create its certificates terminated). RFC 5280 §6.1.1 and §4.1.2.5 put the validity check on the RELYING PARTY and frame the window as a CA warranty, and RFC 8446 §4.4.2.2 imposes no validity rule on a server certificate at all — so a presenter refusing its own leaf is a local availability policy, and it belongs only where it protects something.
Nor is a clock-skew leeway constant needed here (Kubernetes carries CertificateBackdate = 5m for exactly that). Backdating exists where the refusal is terminal; this refusal costs at most one poll interval, because the digests of a refused pair are not recorded and Watch re-examines it.
func (*CertReloader) SetClock ¶ added in v0.12.0
func (r *CertReloader) SetClock(clk clock.Clock)
SetClock overrides the clock the VALIDITY WINDOW is evaluated against, so a harness can express an expired or not-yet-valid certificate as a fixture with fixed dates instead of waiting for real time to pass. A nil clock is ignored. Call it before the first CertReloader.Reload that must observe it; the initial load performed by NewCertReloader never consults the clock, so setting it immediately after construction is in time for every swap.
It governs the validity comparison and NOTHING else. In particular CertReloader.Watch keeps its own time.Ticker on real time deliberately: a poller driven by a frozen fake clock would never tick, and the harness arm that proves onError fires over a broken pair depends on real polling.
Why this is exported, and why it is still not a public API ¶
It is a MODULE-INTERNAL seam, exactly as Server.SetClock is. clock.Clock lives under internal/ and its method set returns internal types (clock.Timer, clock.Ticker), so no package outside this module can name the parameter type or structurally implement it: the method is unreachable from outside GoGraph even though its name is capitalised. An export_test.go wrapper would not do — a _test.go file is compiled only into its own package's test binary, so internal/sim could never reach it.
func (*CertReloader) Watch ¶
func (r *CertReloader) Watch(interval time.Duration, stop <-chan struct{})
Watch starts a background goroutine that calls Reload every interval; Reload itself decides whether the bytes on disk differ from those in service and skips the parse when they do not. The goroutine exits when stop is closed. Watch returns immediately; pair it with sync.WaitGroup if the caller wants to block on shutdown.
Common usage:
stop := make(chan struct{})
go reloader.Watch(30*time.Second, stop)
defer close(stop)
Errors from Reload are surfaced via the onError callback installed at construction time; Watch itself never returns an error.
type Identity ¶
type Identity struct {
// Principal is the authenticated username or identifier.
Principal string
}
Identity carries the authenticated principal's metadata after a successful authentication exchange.
type NegativeTimeoutError ¶ added in v0.14.1
type NegativeTimeoutError struct {
// Field is the Options field name, e.g. "ConnTimeout".
Field string
// Value is the negative duration the caller supplied.
Value time.Duration
}
NegativeTimeoutError reports one Options timeout field that was given a negative duration. The four timeout fields — ConnTimeout, DefaultTxTimeout, DefaultStatementTimeout and MaxTxIdleTime — read 0 as "disabled" and a positive value as "bounded at that value", which leaves a negative duration with no meaning at all. NewServer refuses it rather than coercing it: coercion would hide the caller's mistake and install a bound they never asked for, and the same field's zero value already expresses "no bound" precisely.
This is a construction-time refusal in the style of ErrNoAuthHandler, not a panic: an invalid configuration value is a recoverable condition validated once at the public API boundary.
func (*NegativeTimeoutError) Error ¶ added in v0.14.1
func (e *NegativeTimeoutError) Error() string
Error implements error.
func (*NegativeTimeoutError) Unwrap ¶ added in v0.14.1
func (e *NegativeTimeoutError) Unwrap() error
Unwrap makes errors.Is(err, ErrNegativeTimeout) report true.
type NoAuthHandler ¶
type NoAuthHandler struct{}
NoAuthHandler accepts any credentials without validation. Suitable for development and testing only.
Because it admits every client, NoAuthHandler is never installed by default: the server is secure-by-default. To run without authentication an embedder must opt in explicitly by setting Options.Auth to a NoAuthHandler{} value. The explicit value is itself the opt-in — self-documenting at the call site and impossible to set by accident — and NewServer logs a loud warning when it sees one. Constructing a server with a nil Options.Auth fails closed with ErrNoAuthHandler. Never expose a NoAuthHandler-backed server on an untrusted network.
NoAuthHandler is safe for concurrent use.
func (NoAuthHandler) Authenticate ¶
func (NoAuthHandler) Authenticate(_, principal, _ string) (Identity, error)
Authenticate implements AuthHandler. It always returns an Identity with the given principal and a nil error.
type Options ¶
type Options struct {
// Auth is the authentication handler invoked during HELLO/LOGON. It is
// the security boundary of the server: every client must satisfy it
// before any Cypher statement executes.
//
// Auth must be set; leave it nil and [NewServer] returns
// [ErrNoAuthHandler]. The server is secure-by-default: a nil Auth is NOT
// silently replaced with an open, accept-everyone handler, so a careless
// embedder writing Options{} cannot accidentally expose an
// unauthenticated server. To enforce credentials, set Auth to a real
// [AuthHandler] such as [BasicAuthHandler]. To run without authentication
// (development or testing only) set Auth: [NoAuthHandler]{} explicitly:
// the explicit NoAuthHandler value is itself the opt-in, it is
// self-documenting at the call site, and it is impossible to set by
// accident. In that case [NewServer] still emits a loud warning.
Auth AuthHandler
// Closer, when non-nil, is the store-level teardown owner for the
// durability stack backing this server's engine — typically a
// *[github.com/FlavioCFOliveira/GoGraph/store.DB] bundling the WAL writer
// and the background checkpointer. The server closes it AFTER it has
// drained every active connection, so it runs the one crash-safe teardown
// order (stop the checkpoint goroutine, then close the WAL) only once no
// in-flight transaction can still be writing. Both documented stop
// mechanisms reach that teardown: [Server.Shutdown] closes it on its
// drain-success branch, and [Server.Serve] closes it on its own exit path
// once its connection drain completes (e.g. when the Serve context is
// cancelled). The close is guarded by a [sync.Once] inside the server, so
// the closer's Close runs exactly once regardless of which path wins or
// whether both run; it need not be idempotent itself. Leave it nil for a
// store-less engine or when the embedder tears the durability stack down
// itself; the server then closes nothing beyond its connections.
Closer io.Closer
// TLSConfig, when non-nil, wraps accepted connections with TLS using
// the given configuration verbatim. nil means plain TCP (no TLS).
//
// The server applies no MinVersion or cipher policy of its own: whatever
// config is supplied here is used as-is. To start from a hardened baseline
// (TLS 1.2 floor, modern AEAD/ECDHE cipher list), begin with
// [DefaultTLSConfig] and add your own Certificates or GetCertificate before
// assigning it here.
TLSConfig *tls.Config
// Logger is the structured logger for server events, per-connection and
// per-session alike: accept-loop and handshake events, and every session-level
// event (authentication outcomes, query, BEGIN and COMMIT failures, and the
// transaction-quota refusal). When nil, the default slog handler is used.
//
// The session half was added in rmp #2481: until then newSession hard-coded
// slog.Default(), so an embedder that configured a logger still had the
// majority of the server's events written to the process default.
Logger *slog.Logger
// DatabaseName is the name this server reports as the database serving a
// result, in the `db` field of the RUN and terminal PULL/DISCARD SUCCESS
// metadata. Empty defaults to [DefaultDatabaseName].
//
// GoGraph serves exactly one graph per server, so this is a label rather
// than a selector: a client that names a database in its session config has
// that name echoed back (so its own bookkeeping stays consistent), and a
// client that names none is told this value. The name is not validated and
// an unknown name is not rejected — see [DefaultDatabaseName].
//
// Sending it at all matters because the field is not optional in practice:
// the official neo4j-go-driver returns a nil DatabaseInfo from
// ResultSummary.Database() when `db` is absent, so the idiomatic
// summary.Database().Name() panics with a nil dereference inside the driver
// (rmp #2172).
DatabaseName string
// MaxTxIdleTime bounds how long an OPEN explicit transaction may go without
// the client sending a message, after which it is rolled back and the MVCC read
// snapshot it pinned is released.
//
// Three-way contract, the same for all four timeout fields:
//
// - 0 (the zero value) DISABLES the bound. This is the default — see
// [DefaultMaxTxIdleTime] for why, and for what leaving it disabled costs.
// - A positive value is used verbatim; nothing substitutes a default for it.
// - A NEGATIVE value is an error: [NewServer] refuses it with a
// [*NegativeTimeoutError] rather than coercing it to a default, because
// coercion hides the caller's mistake and installs a bound they never
// asked for.
//
// This is NOT DefaultTxTimeout. That bounds a transaction's total life however
// active it is; this reclaims one that has stopped talking. A busy transaction
// resets it on every message and is limited only by the total bound.
MaxTxIdleTime time.Duration
// MaxOpenTxPerPrincipal caps how many explicit transactions one authenticated
// principal may hold open at once across all of its connections. Exceeding it
// fails the BEGIN with Neo.TransientError.Transaction.MaximumTransactionLimitReached rather than
// queueing. Zero defaults to [DefaultMaxOpenTxPerPrincipal]; a NEGATIVE value
// disables enforcement, which is deliberate and visible at the call site.
MaxOpenTxPerPrincipal int
// MaxConnections is the upper bound on concurrent accepted connections.
// Zero or negative values default to 1024.
MaxConnections int
// MaxMessageBytes caps the cumulative payload size of a single Bolt
// message reassembled from per-chunk fragments. Zero or negative
// values default to [proto.DefaultMaxMessageBytes] (16 MiB).
// Bolt's wire format limits each chunk to 65535 bytes but the
// chunk count is unbounded; this cap closes the Slowloris-style
// DoS vector in which a malicious client streams non-zero chunks
// indefinitely until the server OOMs.
MaxMessageBytes int
// MaxInboundDecodeBytes is the engine-wide ceiling, in bytes, on the total
// decoded-collection memory in flight across ALL connections while messages
// are being decoded. MaxMessageBytes bounds a single message and the decoder
// bounds a single message's decoded collections, but without this aggregate
// bound that per-message cap times MaxConnections is unbounded and reachable
// before authentication — a memory-exhaustion DoS (CWE-770). When the pool is
// drawn down, further inbound decodes fail fast with a retryable transient
// error (backpressure) rather than allocating.
//
// Interpretation mirrors cypher.EngineOptions.GlobalMaxResultBytes:
// - 0 (the zero value) → derive from the Go soft memory limit: one eighth
// of GOMEMLIMIT when the operator has set one (results already claim half,
// and inbound decode is transient), else unlimited. This gives default-on
// protection precisely when a memory budget is declared, and never rejects
// a legitimate workload on a host whose memory the module cannot know.
// - [MaxInboundDecodeBytesUnlimited] (-1) → unlimited (explicit opt-out).
// - a positive value → used verbatim.
//
// GOMEMLIMIT alone cannot mitigate the DoS: in-flight decoded values are live,
// non-collectable memory during decode, so the ceiling must gate before the
// allocation.
//
// Scope: the ceiling bounds the concurrent DECODE-phase memory across
// connections — the simultaneous-allocation vector that is the actual OOM
// risk. A message's reserved bytes are returned as soon as it finishes
// decoding, so this does not bound a message's lifetime while its handler
// runs; a single connection may still transiently hold one decoded message
// (bounded by MaxMessageBytes) during handling. This mirrors the result
// ceiling's transient-vs-lifetime asymmetry.
MaxInboundDecodeBytes int64
// MaxInFlightPerConnection caps the total number of RUN statements
// that may be issued within a single explicit transaction before
// COMMIT or ROLLBACK. Zero or negative values default to
// [DefaultMaxInFlightPerConnection] (1024). The count includes both
// open (not yet fully PULL'd) and already-drained cursors
// accumulated in tx.results since BEGIN; auto-commit cursors are
// not counted (the Bolt v5 state machine already prevents two
// concurrent auto-commit streams). The cap surfaces as a typed
// Bolt FAILURE with code
// "Neo.TransientError.Transaction.MaximumTransactionLimitReached" — TRANSIENT, so a
// driver retries, and the session stays in READY so it can (rmp #2561).
MaxInFlightPerConnection int
// ConnTimeout is the per-connection idle read deadline applied throughout
// the post-handshake message loop. When it is positive, the deadline is reset
// to now+ConnTimeout before every read.
//
// Three-way contract, the same for all four timeout fields:
//
// - 0 (the zero value) DISABLES the deadline: reads carry no deadline at
// all. This is the default — see [DefaultConnTimeout] for why, and for
// what leaving it disabled costs. Liveness against a dead peer comes from
// TCP keep-alive instead (see [DefaultKeepAliveIdle]).
// - A positive value is used verbatim; nothing substitutes a default for it.
// - A NEGATIVE value is an error: [NewServer] refuses it with a
// [*NegativeTimeoutError].
//
// Read the name as an IDLE bound with care: the reader goroutine sits in the
// read while the message loop executes the client's own statement, so a
// positive value also bounds a single long statement — a client waiting for
// the records it asked for is silent, and this deadline cannot tell that from
// abandonment. Size it above the slowest statement the deployment expects,
// not merely above its think-time.
//
// The unauthenticated version-negotiation handshake is bounded separately and
// unconditionally; see [DefaultHandshakeTimeout]. The window between that
// handshake and a successful LOGON is bounded by THIS field alone, so a
// server exposed to an untrusted network wants it set.
ConnTimeout time.Duration
// MaxStatementTimeout is the server-side upper bound on per-statement
// execution time. When a client supplies a timeout via the RUN or BEGIN
// extra metadata, it is silently clamped to MaxStatementTimeout. When
// a client supplies no timeout and MaxStatementTimeout is positive, the
// server applies MaxStatementTimeout unconditionally. Zero means no
// server-side cap (client controls its own timeout).
MaxStatementTimeout time.Duration
// DefaultTxTimeout is the total wall-clock bound applied to an explicit
// transaction (opened by BEGIN) when the client supplies no tx_timeout of its
// own, however busy the transaction is.
//
// Three-way contract, the same for all four timeout fields:
//
// - 0 (the zero value) DISABLES the bound. This is the default — see
// [DefaultTxTimeout] for why, and for what leaving it disabled costs.
// - A positive value is used verbatim; nothing substitutes a default for it.
// - A NEGATIVE value is an error: [NewServer] refuses it with a
// [*NegativeTimeoutError].
//
// A client-supplied tx_timeout takes precedence; MaxStatementTimeout, when
// set, additionally clamps the effective value.
DefaultTxTimeout time.Duration
// DefaultStatementTimeout is the wall-clock bound applied to an AUTOCOMMIT
// statement (a bare RUN outside an explicit transaction) when the client
// supplies no per-statement `timeout` of its own. It is the autocommit
// counterpart of DefaultTxTimeout.
//
// Three-way contract, the same for all four timeout fields:
//
// - 0 (the zero value) DISABLES the bound. This is the default — see
// [DefaultStatementTimeout] for why, and for what leaving it disabled
// costs (#1828 is the defect report that argued for a finite floor here;
// read it before deciding a deployment can do without one).
// - A positive value is used verbatim; nothing substitutes a default for it.
// - A NEGATIVE value is an error: [NewServer] refuses it with a
// [*NegativeTimeoutError].
//
// A client-supplied `timeout` takes precedence; MaxStatementTimeout, when set,
// additionally clamps the effective value.
DefaultStatementTimeout time.Duration
}
Options configures a Server. It is a plain configuration value read once by NewServer; it is safe for concurrent read use once constructed, but must not be mutated after being passed to NewServer. The referenced TLSConfig, Auth, Logger, and Closer carry their own concurrency contracts.
type Server ¶
type Server struct {
// contains filtered or unexported fields
}
Server is the Bolt v5 TCP server. It accepts connections from a net.Listener, negotiates the protocol version, and runs the Bolt message loop on each connection.
Server is safe for concurrent use by multiple goroutines.
func NewServer ¶
NewServer creates a Server backed by eng. Zero-value Options fields are filled with sensible defaults.
NewServer is secure-by-default: it never silently installs an accept-everyone authentication handler. If Options.Auth is nil it fails closed and returns ErrNoAuthHandler so that an unauthenticated server is never started by accident. To run without authentication on purpose (development or testing), set Options.Auth to a NoAuthHandler{} value explicitly: NewServer then admits every client and logs a loud warning that the operator has knowingly disabled authentication. The explicit NoAuthHandler value is itself the opt-in — self-documenting at the call site and impossible to set by accident. When Options.Auth is any other (real) handler it is used as-is.
func (*Server) ListenAndServe ¶
ListenAndServe creates a TCP listener on addr and calls Serve. It blocks until the server stops. The listener is closed when Serve returns.
func (*Server) Serve ¶
Serve accepts connections from ln until ctx is cancelled or Shutdown is called. It blocks until all active connections have closed. The provided ln is closed by Serve when the accept loop exits.
Once every connection has drained, Serve also closes the owned Options.Closer (the store-level teardown owner for the durability stack, typically a *github.com/FlavioCFOliveira/GoGraph/store.DB), so stopping the server by cancelling ctx tears the WAL/checkpoint stack down in its crash-safe order exactly as Server.Shutdown does — no checkpoint goroutine or WAL handle outlives Serve. The close happens strictly after the drain, so it can never race an in-flight write, and it is once-guarded, so a subsequent (or concurrent) Shutdown does not close the closer again. A failed close is returned (joined with any accept error) rather than swallowed. After Serve returns, the durability stack is closed: the Server must not be reused to serve writes again.
Example ¶
ExampleServer_Serve starts a Bolt server backed by an in-memory graph, connects a Bolt client, and runs a query over the session. The listener binds to 127.0.0.1:0 so the OS assigns a free port. Teardown closes the client first, then cancels Serve and waits for it to drain every connection goroutine — leaving no leaked goroutine behind.
The network round-trip is non-deterministic in timing, so the example asserts the deterministic query result rather than any wire-level output.
package main
import (
"context"
"fmt"
"net"
"time"
"github.com/FlavioCFOliveira/GoGraph/bolt/server"
"github.com/FlavioCFOliveira/GoGraph/cypher"
"github.com/FlavioCFOliveira/GoGraph/graph/adjlist"
"github.com/FlavioCFOliveira/GoGraph/graph/lpg"
"github.com/neo4j/neo4j-go-driver/v5/neo4j"
)
func main() {
// Engine over an empty in-memory labelled property graph.
g := lpg.New[string, float64](adjlist.Config{})
eng := cypher.NewEngine(g)
// The explicit NoAuthHandler{} value is the opt-in that lets this example
// run without credentials; the server is secure-by-default and otherwise
// refuses to start with a nil Auth handler.
srv, err := server.NewServer(eng, server.Options{ConnTimeout: 5 * time.Second, Auth: server.NoAuthHandler{}})
if err != nil {
fmt.Println("new server:", err)
return
}
// Ephemeral port; ln.Addr() reveals the chosen port for the client.
ln, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
fmt.Println("listen:", err)
return
}
addr := ln.Addr().String()
ctx, cancel := context.WithCancel(context.Background())
serveErr := make(chan error, 1)
go func() { serveErr <- srv.Serve(ctx, ln) }()
// Connect a Bolt client and run a trivial read query.
driver, err := neo4j.NewDriverWithContext("bolt://"+addr, neo4j.NoAuth())
if err != nil {
fmt.Println("driver:", err)
cancel()
<-serveErr
return
}
sess := driver.NewSession(ctx, neo4j.SessionConfig{})
result, err := sess.Run(ctx, "RETURN 1 AS n", nil)
if err != nil {
fmt.Println("run:", err)
} else if rec, err := result.Single(ctx); err != nil {
fmt.Println("single:", err)
} else {
n, _ := rec.Get("n")
fmt.Println("n =", n)
}
_ = sess.Close(ctx)
// Clean shutdown: close the client so server-side connection goroutines
// observe EOF, then cancel Serve and wait for it to return. Serve only
// returns after every connection goroutine has finished.
_ = driver.Close(ctx)
cancel()
<-serveErr
}
Output: n = 1
func (*Server) SetClock ¶ added in v0.12.0
SetClock overrides the server's wall-clock source for the explicit-transaction timeout reaper and for the transaction registry's start and elapsed times, so a harness can drive both on virtual time. A nil clock is ignored. It must be called BEFORE Server.Serve starts accepting connections, because each connection captures the clock at session construction.
Why this is exported, and why it is still not a public API ¶
It is a MODULE-INTERNAL seam. clock.Clock lives under internal/ and its method set returns internal types (clock.Timer, clock.Ticker), so no package outside this module can name the parameter type or structurally implement it: the method is unreachable from outside GoGraph even though its name is capitalised. The alternative — an Options.Clock field — would place an unusable internal type in the public configuration struct of every embedder.
It was previously unexported, reachable only through an export_test.go wrapper whose godoc claimed the DST harness used it. That claim could not have been true: a _test.go file is compiled only into its own package's test binary, so internal/sim could never reach it. rmp #2482 needs the seam for real — the idle-transaction reaper on a fake clock — and this is it.
func (*Server) Shutdown ¶
Shutdown gracefully stops accepting new connections and waits for active connections to finish. If connections do not finish within 30 seconds, it closes the listener forcefully and returns an error.
When the server was constructed with Options.Closer (the store-level teardown owner for the durability stack, typically a *github.com/FlavioCFOliveira/GoGraph/store.DB), Shutdown closes it AFTER every active connection has drained — so the WAL/checkpoint teardown runs in its crash-safe order only once no in-flight transaction can still be writing. Closing it before the drain could let a still-executing write race the WAL close. The closer is therefore NOT torn down on the timeout or ctx-cancellation paths: an undrained connection may still hold a transaction, so tearing the WAL down underneath it is exactly what must be avoided; in those cases the connections are abandoned. A still-running Server.Serve remains blocked on the same drain, and when the abandoned connections do eventually finish (idle timeout, transaction reap, client exit), Serve's own exit path performs the post-drain close — so the closer is torn down as soon as a full drain truly completes, and is left for process exit only if it never does. The close is once-guarded: whichever of Serve or Shutdown drains first runs it, and the other observes the same cached result, so the closer is never closed twice (including on a double Shutdown). A failed WAL close is surfaced rather than swallowed.
func (*Server) TerminateTransaction ¶ added in v0.11.0
TerminateTransaction rolls back the open transaction with the given id. It returns ErrNoSuchTransaction if no such transaction is open.
It releases NO writer serialisation and no visibility barrier — rmp #2305/#2306 retired both, exactly as Server.Transactions above says at length. This line claimed otherwise until rmp #2560, contradicting its own neighbour twenty lines up. What the rollback reclaims is the transaction's unpublished commit record and its reclamation-horizon slot.
The terminated connection is told so on its next request-phase message: [Session.terminateTxByOperator] arms Neo.ClientError.Transaction.Terminated, which is deliberately distinct from the code an expired bound produces, so an operator termination is not reported to the client as a timeout that never happened.
The rollback is performed by the connection that owns the transaction, on its own goroutine, because a Session is single-threaded by contract. This call therefore REQUESTS the rollback and returns once the request is delivered; the transaction's context is cancelled synchronously, so a statement already executing is interrupted immediately, and the rollback follows as soon as the owning loop observes the request. Use Server.Transactions to confirm it has gone.
The rollback is atomic: it unwinds every statement of the transaction, exactly as a client ROLLBACK would, so no partial state is left behind.
TerminateTransaction is safe to call from any goroutine.
func (*Server) Transactions ¶ added in v0.11.0
func (s *Server) Transactions() []TransactionInfo
Transactions returns a snapshot of every explicit transaction currently open on this server, oldest first.
It is the diagnostic half of the pair Server.TerminateTransaction completes.
The reason it matters CHANGED with rmp #2305/#2306, and the old reason is worth stating so nobody restores it: an open writing transaction used to hold the engine's writer serialisation and the visibility barrier for its whole lifetime, so "while one is open every reader waits" was literally true and an abandoned transaction was an outage. It no longer holds either. What an abandoned transaction pins now is the reclamation horizon — no version it could still read is freed while it lives — so the symptom is unbounded version memory rather than stalled clients, and finding it still requires knowing its principal, its age and its current statement.
Concurrency and consistency contract ¶
Transactions is safe to call from any goroutine at any time, including while the server is serving. The returned slice and its elements are copies, and nothing the caller does with them can affect the server.
A listing is NOT a single global instant, and that is a deliberate trade rather than an oversight. It is consistent PER ENTRY and approximate ACROSS entries:
- Within one TransactionInfo every field belongs to one instant. ID, Principal, Remote, Mode and StartedAt are fixed when the transaction opens and never move; State and Query are read together with one atomic load, so the pair is always one the transaction genuinely passed through.
- Across two entries the instants may differ by however long the listing takes to walk from one to the other. A listing can therefore show a combination of per-transaction states that never held simultaneously across the whole server. The SET of transactions listed is exact — membership is read under the registry's lock — so an entry is never invented, duplicated or lost; only the moving fields of different entries may be skewed.
What that bought: refreshing a transaction's state used to take a process-global mutex after EVERY inbound Bolt message, which measured 8.3 ns at one goroutine and 77.5 ns at eight — the registry's most frequent operation was also its worst-scaling one (rmp #2714). It now takes no lock.
The weaker guarantee is sound for what this listing is for. Deciding that a transaction is abandoned, or long-running, or the one to terminate, is a judgement about ONE transaction, and each entry is internally exact. Do NOT build a cross-transaction invariant on a listing — "these two were in TX_STREAMING at the same moment" is not something a listing can establish, and it could not establish it before this change either, because the transactions were free to move on the instant the lock was released.
type Session ¶
type Session struct {
// contains filtered or unexported fields
}
Session holds all per-connection state for a single Bolt v5 client connection.
Session is NOT safe for concurrent use. Each accepted TCP connection owns exactly one Session, and the message loop is single-threaded per connection.
func (*Session) Close ¶ added in v0.2.0
func (s *Session) Close()
Close tears the session down on connection teardown: it drains any open cursor and rolls back any open explicit transaction so the engine writer serialisation is released immediately rather than lingering until the GC finalises the leaked Result/transaction (#1309). It is safe to call exactly once from the connection handler's deferred cleanup on every exit path (clean close, read or write error, panic). Idempotent: a second call, or a call on a session with no open transaction, is a no-op.
An explicit transaction still open at this point is an abnormal disconnect — the client dropped the connection (or hit an idle timeout, or the handler panicked) without sending COMMIT, ROLLBACK, or RESET. Close counts it as [metricTxAbandoned] before the rollback so an operator can distinguish a leaked transaction reclaimed here from one ended in an orderly way. This is the only site that emits tx.abandoned: a FAILED-transition reclaim (#1312) goes through [Session.abortTx] directly and is an in-session state change, not a disconnect, so it is not counted abandoned.
func (*Session) HandleMessage ¶
HandleMessage dispatches msg to the correct per-state handler and returns the response messages to send to the client.
On an illegal state transition or internal error the session moves to FAILED and HandleMessage returns a single *proto.Failure response. The caller is responsible for encoding and sending all returned messages.
When a record sink is installed (see [Session.setRecordSink]), the RECORD messages of a PULL are written through the sink as the cursor is iterated and are NOT part of the returned slice, which then carries only the trailing SUCCESS or FAILURE. A sink write failure is surfaced as an error wrapping [errRecordWrite]: the connection framing is unrecoverable and the caller must tear the connection down without writing anything further.
Observability ¶
Every call records a latency observation under "bolt.server.HandleMessage.message.<type>", the module's per-message Bolt histogram; see msgmetrics.go for the naming and for why the message type is carried in the name. The window is dispatch plus handler execution — it EXCLUDES the framing read that produced msg and the trailing response write the caller performs, and INCLUDES the RECORD writes a PULL streams through its own sink. It is emitted here, on the exported entry point, rather than in the serve loop, so a session driven directly through HandleMessage is observed on exactly the same terms as one driven over a socket.
type State ¶
type State uint8
State represents the Bolt v5 per-connection protocol state machine state.
const ( // StateConnected is the initial state: TCP connection established, no // protocol negotiation has occurred yet. StateConnected State = iota // StateNegotiation is reached after version negotiation; the server awaits // the client's HELLO message. StateNegotiation // StateAuthentication is the pre-LOGON state reached after a successful // credential-less HELLO on Bolt >= 5.1. Bolt 5.1 split authentication out of // HELLO into a dedicated LOGON message, so a 5.1+ client sends a HELLO // carrying only driver metadata and then a LOGON carrying the credentials. // In this state the connection is not yet authenticated; only LOGON, LOGOFF, // RESET, and GOODBYE are legal, and a successful LOGON transitions to // StateReady. On Bolt <= 5.0 (and the white-box tests, which run at the // zero-value version) HELLO authenticates inline and goes straight to // StateReady, so this state is never entered. (task #1470) StateAuthentication // StateReady is the idle state after a successful HELLO or after a result // set has been fully consumed, committed, or rolled back. StateReady // StateStreaming is active when a query has been run (auto-commit) and // records are available to pull. StateStreaming // StateTxReady is reached after BEGIN; the server awaits RUN, COMMIT, or // ROLLBACK within an explicit transaction. StateTxReady // StateTxStreaming is active when a query has been run inside an explicit // transaction and records are available to pull. StateTxStreaming // StateFailed is entered when a request fails; the server ignores further // requests until RESET is received. StateFailed // StateDefunct is the terminal state: the connection is closed and no // further messages are processed. StateDefunct )
func HelloTransition ¶ added in v0.3.0
HelloTransition computes the next state for a successful HELLO given the negotiated Bolt version. It is the version-aware variant of the NEGOTIATION→HELLO branch of Transition:
- Bolt <= 5.0 (and the zero-value version used by direct white-box tests): HELLO authenticates inline and advances straight to StateReady.
- Bolt >= 5.1: HELLO is credential-less by spec, so a successful HELLO advances to the pre-LOGON StateAuthentication, from which a successful LOGON reaches StateReady.
It is only valid in StateNegotiation; any other current state, or a failed HELLO, is delegated to Transition, which returns StateFailed (with ErrInvalidTransition for an illegal current state). (task #1470)
func StreamingTransition ¶
StreamingTransition is a variant of Transition for PULL in STREAMING or TX_STREAMING states when there are more records to deliver (has_more=true). In that case the connection remains in the same streaming state instead of returning to READY/TX_READY.
func Transition ¶
Transition computes the next state given the current state, the incoming message, and whether the operation succeeded.
msg must be one of the pointer types from the proto package (e.g. *proto.Run, *proto.Pull, etc.). success indicates whether the server-side operation succeeded; on failure the next state is StateFailed (unless the transition itself is illegal).
Returns (StateFailed, ErrInvalidTransition) for illegal state/message combinations.
The one documented exception: a refused RESOURCE CAP (rmp #2561) ¶
"On failure the next state is StateFailed" describes a failure to CARRY OUT the request. A request refused because a bounded resource is momentarily full is not that: [Session.handleBegin]'s open-transaction-quota branch returns its FAILURE without calling Transition at all, and the session stays in READY.
That is deliberate. A cap is back-pressure — the slot frees when another of the principal's transactions closes — so the right response is to retry the same BEGIN, and requiring a RESET first would charge the client a round trip for the server being busy. The refusal carries a TRANSIENT code for the same reason.
The adjacent newTx failure in the same function DOES go to StateFailed, because that is a genuine failure to open a transaction. The two differ on purpose; until rmp #2561 they differed by accident and nothing said so.
Bolt v5 state machine has O(states×messages) branches; splitting it would obscure the protocol spec.
type TransactionInfo ¶ added in v0.11.0
type TransactionInfo struct {
// StartedAt is when the BEGIN completed, from the server's clock.
StartedAt time.Time
// ID identifies the transaction for [Server.TerminateTransaction]. It is
// unique for the lifetime of the server: a new transaction on the same
// connection gets a new ID, so an ID captured from an earlier listing can
// never terminate a later transaction that happens to reuse the connection.
ID string
// Principal is the authenticated identity that opened the transaction, empty
// when the server runs without authentication and the client sent no
// principal.
Principal string
// Remote is the client's network address, as the server sees it.
Remote string
// Mode is "w" for a writing transaction or "r" for a read-only one.
//
// NEITHER blocks anybody as of rmp #2305/#2306. A writing transaction used to
// hold the engine's writer serialisation and the visibility barrier for its whole
// lifetime, so an abandoned one was an outage and finding it was urgent; it now
// holds only its own unpublished commit record and a reclamation-horizon slot.
// What an abandoned writing transaction still costs is therefore version memory —
// no version it can reach is reclaimable while it lives — not other clients'
// progress.
Mode string
// State is the Bolt state machine state of the owning session, rendered for a
// human — "TX_READY" between statements, "TX_STREAMING" while a result is
// being drained.
State string
// Query is the text of the most recent statement RUN inside the transaction,
// empty when BEGIN has not yet been followed by a RUN. It is the field that
// answers "what is it doing?", which a counter and a log line cannot.
Query string
// Elapsed is how long the transaction has been open, measured at snapshot
// time. It is provided rather than left to the caller because the server's
// clock may be injected, and subtracting StartedAt from time.Now() would then
// be wrong.
Elapsed time.Duration
}
TransactionInfo describes one explicit Bolt transaction that is currently open. The values are copies, so holding one blocks nothing and observes no later change.
Every field of ONE TransactionInfo describes one instant: Server.Transactions takes ID, Principal, Remote, Mode and StartedAt from fields that never change after the transaction opens, and takes State and Query together with a single atomic load. A State and a Query that appear side by side here really were simultaneously true. Two DIFFERENT entries of the same listing need not share an instant — see the consistency note on Server.Transactions.
TransactionInfo is safe for concurrent use because it is immutable once returned by Server.Transactions.
type Tx ¶
type Tx struct {
// contains filtered or unexported fields
}
Tx wraps an engine-level explicit transaction (cypher.ExplicitTx) for a single Bolt transaction opened by a BEGIN message. Every RUN issued between BEGIN and COMMIT/ROLLBACK executes against the SAME underlying engine transaction, so the statements are atomic together: COMMIT makes them durable and visible as one unit, ROLLBACK unwinds all of them (#1280). This replaces the previous behaviour in which each RUN opened and committed its own autocommit transaction and ROLLBACK undid nothing.
Tx is NOT safe for concurrent use; it is owned by a single Session whose message loop is single-threaded per connection.
func (*Tx) Commit ¶
Commit makes every statement issued since BEGIN durable and visible as one atomic unit, then releases the transaction's resources. The engine fsyncs the WAL exactly once for the whole transaction (WAL-backed) and commits the secondary-index buffer; on a store-less engine the writes are already visible and the index buffer is finalised. The writer serialisation is released.
Open result cursors are closed first (releasing their iterator state); the commit decision itself is made by the engine transaction, not by the cursors.
func (*Tx) Rollback ¶
Rollback unwinds every statement issued since BEGIN — restoring the in-memory graph to its pre-transaction state via the engine's accumulated undo log, and (WAL-backed) discarding the WAL transaction so a fresh recovery observes none of the writes — then releases the transaction's resources. It is best-effort and always releases the writer serialisation, even if an inverse operation fails.
func (*Tx) Run ¶
Run executes query inside the transaction WITHOUT committing, buffers the result cursor, and returns it to the caller for streaming. The statement's writes accumulate in the engine transaction and become durable/visible only on Commit.
Runtime pipeline errors (where the query compiled and executed under the visibility barrier but the execution pipeline failed — e.g. a constraint violation or a type error mid-pipeline) are wrapped in a zero-row cypher.Result whose Err() carries the error. This preserves the Bolt v5 state-machine contract: the session stays in TX_STREAMING so the driver can drain the cursor via PULL (where the FAILURE surfaces), rather than receiving a FAILURE directly from RUN. Build-phase errors (context cancellation, parse/sema/plan failures, DDL rejection) are propagated as a non-nil error return so the server enters FAILED at RUN time.