Documentation
¶
Overview ¶
Package executor is a Go client for the executor door of a running flexiq-server: it attaches, runs the jobs it is given, and reports what happened.
It is the other half of github.com/ByteVeda/flexiq/sdks/go/v2, which submits work. The two are separate packages because they are separate doors behind separate scopes — a token carrying "execute" cannot enqueue, and a token carrying "produce" cannot attach.
An executor cannot enqueue ¶
There is no enqueue-shaped RPC in flexiq.executor.v1 at all, so this is not a restriction this package applies — it is one the wire has. A task that fans out to a second stage goes back through the producer door as an ordinary client, holding a second, produce-scoped credential of its own.
Getting started ¶
w, err := executor.New("queue.internal:50051",
executor.WithToken(os.Getenv("FLEXIQ_EXECUTE_TOKEN")),
executor.WithID("go-worker-1"),
executor.WithSlots(8),
)
if err != nil {
return err
}
defer w.Close()
w.Handle("billing.charge", func(ctx context.Context, job *executor.Job) (any, error) {
var charge Charge
if err := job.Bind(&charge); err != nil {
return nil, executor.Fatal(err)
}
return Receipt{ID: charge.OrderID}, nil
})
err = w.Run(ctx)
Worker.Run blocks. It attaches, dispatches jobs to their handlers, and reconnects when the scheduler rotates the stream — which it does every half hour by default, because a gRPC stream cannot be load balanced once it has started. A rotation is not a failure and does not cost a job in flight.
What a handler returns ¶
A value and no error is a success; a nil value means the task returned nothing, which is a different answer from returning an empty one. An error is a failure the scheduler will retry, unless it is wrapped in Fatal. Whether to retry is the executor's decision: only it can see the exception, and the scheduler never inspects one.
A handler's context carries the job's timeout and is cancelled when the scheduler asks for the job to stop. A handler that returns its context's error after a cancel settles the job as cancelled rather than failed.
Durable steps ¶
A step runs once per job, not once per attempt. Its result is recorded, and a later attempt replays that record instead of running the body again — which is what stops a retry charging a card twice.
w.Handle("billing.charge", func(ctx context.Context, job *executor.Job) (any, error) {
var order Order
if err := job.Bind(&order); err != nil {
return nil, executor.Fatal(err)
}
receipt, err := executor.Step(ctx, job, "charge",
func(ctx context.Context, key string) (Receipt, error) {
return gateway.Charge(ctx, order, key)
})
if err != nil {
return nil, err
}
if err := job.Sleep(ctx, "settlement", 24*time.Hour); err != nil {
return nil, err
}
return receipt, nil
})
Step is a package function rather than a method because a Go method cannot have a type parameter. StepKeyed identifies a step by data instead of by position, and Job.Sleep ends the attempt and reschedules the job.
The key handed to the body is this step's downstream idempotency key, stable across every attempt. Memoization closes the replay window; only a key the other service dedupes on closes the crash window between a remote call succeeding and its row committing.
Unlike the side channel, this capability **fails rather than degrades**: a step call against a scheduler that does not offer a step store returns a retryable StepError rather than running un-memoized. There is no version of "your charge step silently lost its memo" that beats a failure naming the reason.
A refusal is the task body's to handle: catch it, do something else, return a value, and the job is recorded a success. Two are not, and returning normally past either does not make them so. ErrStepDiverged fails the attempt anyway, because the deployed code and the recorded rows disagree. ErrStepSuperseded settles nothing at all — another attempt owns the job, and this one may not write over it.
What this package does not do ¶
Task registration, middleware, the admin surface, settings and pub/sub are absent because they are not on this door — see contracts/REMOTE_SDK_CONTRACT.md.
Index ¶
- Constants
- Variables
- func Fatal(err error) error
- func Step[T any](ctx context.Context, job *Job, name string, ...) (T, error)
- func StepKeyed[T any](ctx context.Context, job *Job, name, key string, ...) (T, error)
- type AttachError
- type Handler
- type Job
- func (j *Job) Bind(targets ...any) error
- func (j *Job) Call() (flexiq.Call, error)
- func (j *Job) Log(level, message string, extra any) error
- func (j *Job) Payload() []byte
- func (j *Job) Progress(percent int)
- func (j *Job) Publish(value any) error
- func (j *Job) RunKey() string
- func (j *Job) Sleep(ctx context.Context, name string, d time.Duration) error
- func (j *Job) SleepUntil(ctx context.Context, name string, at time.Time) error
- type Option
- func WithGRPCDialOptions(opts ...grpc.DialOption) Option
- func WithHandshakeTimeout(d time.Duration) Option
- func WithHeartbeatInterval(d time.Duration) Option
- func WithID(id string) Option
- func WithInsecureTransport() Option
- func WithLogger(logger *slog.Logger) Option
- func WithMaxMessageBytes(n int) Option
- func WithReconnectBackoff(minDelay, maxDelay time.Duration) Option
- func WithSDK(name, version string) Option
- func WithShutdownDrain(d time.Duration) Option
- func WithSlots(slots int) Option
- func WithStepAckTimeout(d time.Duration) Option
- func WithStepLimits(limits StepLimits) Option
- func WithTLS(cfg *tls.Config) Option
- func WithToken(token string) Option
- func WithTransportCredentials(creds credentials.TransportCredentials) Option
- func WithUserAgent(agent string) Option
- type StepError
- type StepLimits
- type Worker
Constants ¶
const ( // LevelInfo is an ordinary log line. LevelInfo = "info" // LevelWarn is a log line worth an operator's attention. LevelWarn = "warn" // LevelError is a log line describing something that went wrong. LevelError = "error" // LevelResult marks a published partial: the message is empty and the value // rides in the frame's extra. [Job.Publish] writes one. LevelResult = "result" )
Log levels a task log frame carries. Any string is accepted by the scheduler; these are the ones the rest of the system reads.
const ( // CapSideChannel unlocks progress and log frames. Absent, they are no-ops. CapSideChannel = "side_channel" // CapLease unlocks a lease on every dispatch, echoed on every frame about // it. Absent, frames carry none. CapLease = "lease" // CapSteps unlocks durable steps: the scheduler applies step commits and // sends each dispatch its recorded snapshot. Unlike the other two this one // **fails rather than degrades** — a durable step that silently did not // commit is a step that will re-run a charge — so a step call against a // scheduler that did not acknowledge it fails retryably instead. CapSteps = "steps" )
The optional behaviours an executor can take part in. A client sends no frame for a behaviour it did not advertise, and the scheduler sends none for one it did not acknowledge.
const ( // DefaultMaxStepBytes is the largest encoded result one step may commit: a // checkpoint, not a data payload. DefaultMaxStepBytes = step.DefaultMaxStepBytes // DefaultMaxTotalBytes is the largest total across every committed step of // one job. The snapshot is loaded whole at attempt start, which a per-step // cap alone does not bound once a loop runs ten thousand times. DefaultMaxTotalBytes = step.DefaultMaxTotalBytes // DefaultMaxSteps caps how many steps one job may commit. A loop of cheap // steps returning nothing slips past a byte cap. DefaultMaxSteps = step.DefaultMaxSteps )
Caps a job's committed steps are held to, by default.
The scheduler's own numbers, copied rather than invented. Checking them here buys the good error message — one that names the step and the value the caller passed — while the check that actually holds is the scheduler's, inside the write's own transaction. Configure them only to match a scheduler configured away from these.
const MaxMessageBytes = 68 * 1024 * 1024
MaxMessageBytes is the executor door's per-message cap: the 64 MiB the worker frame protocol allows a payload, plus 4 MiB of envelope headroom.
It is deliberately not the producer door's 4 MiB. The two doors carry different things, and grpc-go caps receiving at 4 MiB while leaving sending unbounded — so a client that leaves either alone attaches cleanly and fails on its first large job, which is the worst time to find out.
const ProtocolVersion uint32 = 1
ProtocolVersion is the worker frame format this client speaks. Both ends announce it in the handshake and both reject a mismatch; it is not the package version and not the contract level.
Variables ¶
var ( // ErrStepRetryable means the backend failed, not the request. The attempt // fails and the job's retry policy has it. A commit that was never // acknowledged is one of these: nothing confirmed the write landed, so a // replay is safe. ErrStepRetryable = errors.New("flexiq: the step commit failed and the attempt should retry") // ErrStepPermanent means the commit will never succeed — a divergence, a // cap, a bad encoding. Retrying only wastes an attempt, so it dead-letters. ErrStepPermanent = errors.New("flexiq: the step commit will never succeed") // ErrStepSuperseded means another attempt owns this job now. This one must // stop **without writing**: no result frame is sent at all, because the job // is proceeding correctly somewhere else. ErrStepSuperseded = errors.New("flexiq: another attempt owns this job") // a fleet mid-rollout can still place the next attempt somewhere that // commits. It is never a silent un-memoized run — there is no version of // "your charge step lost its memo" that beats a failure naming the reason. ErrStepUnavailable = errors.New("flexiq: the scheduler offers no step store") )
The verdicts a refused step carries, matched with errors.Is.
A refusal is classified by the side that holds storage, because only it can see the error — so the verdict crosses the wire as an enum and the message crosses as prose. Never read the verdict back out of the message.
var ErrCancelled = errors.New("flexiq: job cancelled")
ErrCancelled settles a job as cancelled rather than failed.
A handler returns it, or its context's error, after observing a cancel. A handler that ignores the cancel and finishes normally settles normally — cancellation is cooperative, and nothing here stops a running goroutine.
var ErrFatal = errors.New("flexiq: fatal task error")
ErrFatal marks a task failure the scheduler must not retry.
Match it with errors.Is; produce it with Fatal. Retrying is the default because most failures are transient and the scheduler never inspects an error — only the executor can see the exception, so only the executor can say that running the task again would fail the same way.
var ErrNoToken = errors.New("flexiq: no token: every call to this door carries one, use WithToken")
ErrNoToken is returned by New when no credential was supplied. There is no anonymous path on this door.
var ErrStepDiverged = errors.New("flexiq: the step sequence changed between attempts")
ErrStepDiverged means the running code asked for a different step than the one recorded at that position: the step sequence changed between attempts.
Permanent, and one of the two refusals a task body **cannot** carry on past — ErrStepSuperseded is the other, and settles nothing at all. Every refusal besides those two is the body's to handle however it likes: catch it, do something else, return a value, and this client takes that at its word.
This one says the deployed code and the recorded rows disagree, so a memoized result would answer a different question than the step asking for it, and an attempt that continued would write into a sequence that no longer lines up. Returning normally past it fails the attempt anyway.
The other SDKs enforce that with an exception tier `catch` cannot reach — `BaseException` in Python, `java.lang.Error` in Java. Go has no such tier, so this client checks once, when the handler returns.
Drain or dead-letter a task's in-flight jobs before deploying a change to its step sequence.
var ErrStepSlept = errors.New("flexiq: the attempt ended in a durable sleep")
ErrStepSlept ends the attempt in a durable sleep.
Returned by Job.Sleep and Job.SleepUntil once the deadline is committed: the row is written, the claim is released and the job is already scheduled to wake. Propagate it. Anything the task body does past this point runs unclaimed and will run again when the job wakes.
Swallowing it does not resume the attempt — the claim is gone either way, and this client writes the slept frame regardless.
Functions ¶
func Fatal ¶
Fatal wraps an error so the job dead-letters instead of retrying.
return nil, executor.Fatal(fmt.Errorf("no such customer: %q", id))
The wrapped error keeps its type and message: the type becomes the errtype in the canonical failure JSON, and errors.Is and errors.As still reach it.
func Step ¶
func Step[T any](ctx context.Context, job *Job, name string, body func(ctx context.Context, key string) (T, error), ) (T, error)
Step runs body once for this job and memoizes what it returned.
receipt, err := executor.Step(ctx, job, "charge",
func(ctx context.Context, key string) (Receipt, error) {
return gateway.Charge(ctx, order, key)
})
On a later attempt of the same job the body does **not** run: its recorded result is decoded and returned. That is the whole point — re-running it is the double charge durable steps exist to prevent.
key is the downstream idempotency key for this step, "{run}:{name}#{n}", and it is the same string on every attempt. Memoization closes the replay window, not the crash window: between "the charge succeeded" and "the step row committed" the process can die, and only a key the other service dedupes on closes that. Hand this one to any API that takes one.
A package function rather than a method because a Go method cannot have a type parameter, and a step whose result decodes into your own type is worth more than one that hands back `any`.
Steps are numbered by occurrence within an attempt, so they must be asked for in the same order every time. A loop over anything unordered wants StepKeyed.
func StepKeyed ¶
func StepKeyed[T any](ctx context.Context, job *Job, name, key string, body func(ctx context.Context, key string) (T, error), ) (T, error)
StepKeyed runs body once per key rather than once per position.
total, err := executor.StepKeyed(ctx, job, "refund", order.ID,
func(ctx context.Context, key string) (int64, error) { ... })
A keyed step is matched by its key wherever it sits in the recorded sequence, so a loop over a map — or anything else whose order is not guaranteed — can hand its steps back in a different order without every one of them looking like a different question. It never spends an occurrence, so adding one cannot shift the number of a later unkeyed step.
Each key must be used at most once per attempt.
Types ¶
type AttachError ¶
type AttachError struct {
// Code is the gRPC status code the scheduler refused with.
Code codes.Code
// Message is the scheduler's own text, written for humans.
Message string
// Permanent reports whether reconnecting could ever succeed.
Permanent bool
}
AttachError is a refused attach.
Some refusals are worth retrying and some are not, and the difference is the status code rather than the message. Worker.Run returns a permanent one and reconnects through the rest.
func (*AttachError) GRPCStatus ¶
func (e *AttachError) GRPCStatus() *status.Status
GRPCStatus keeps status.Code and status.Convert working on this error.
type Handler ¶
Handler runs one job.
The value it returns becomes the job's result; a nil value means the task returned nothing, which is a different answer from returning an empty one. An error fails the job and the scheduler retries it, unless it is wrapped in Fatal.
The context carries the job's timeout and is cancelled when the scheduler asks for the job to stop. Honour it: cancellation is cooperative and nothing here can stop a goroutine that does not return.
type Job ¶
type Job struct {
// ID is the job's identifier. It is not permanent: retention archives and
// then deletes, so an id that reads NOT_FOUND later did still exist.
ID string
// TaskName is the task this attempt runs.
TaskName string
// Queue is the queue the job was drawn from.
Queue string
// Namespace the job belongs to, told to the executor and never accepted
// from one. Empty when the server sent none.
Namespace string
// RetryCount is how many attempts already failed.
RetryCount int
// MaxRetries is how many this job is allowed.
MaxRetries int
// Timeout is how long this attempt may run. Zero means no limit.
Timeout time.Duration
// DisabledMiddleware names the middleware an operator has turned off for
// this task, resolved at dispatch. It rides the dispatch because an
// executor has no settings store to read it from.
DisabledMiddleware []string
// Metadata is the job's metadata blob as stored: opaque JSON text,
// byte-preserved. It has no schema this package knows; parse it only if
// your own code wrote it.
Metadata string
// contains filtered or unexported fields
}
Job is one dispatched attempt.
Everything on it was resolved by the scheduler from the dispatch it recorded. Nothing a client sends names a namespace, an owner, an attempt or a resource cap — anything a client could name is something a client could forge.
func (*Job) Bind ¶
Bind decodes the call's positional arguments into targets, in order.
var charge Charge
if err := job.Bind(&charge); err != nil { ... }
The cross-SDK convention is a single object argument, so one target is the usual case. A call carrying keyword arguments is refused rather than bound by position — read those with Job.Call.
func (*Job) Call ¶
Call decodes the whole payload: the positional and keyword arguments the task was enqueued with.
Values arrive as `any`, so a number is whatever CBOR's widest form for it is. Use Job.Bind to decode into your own types instead.
func (*Job) Log ¶
Log records one structured line against this attempt.
extra is marshalled as JSON and may be nil. Fire and forget, and a no-op without the side-channel capability, exactly like Job.Progress.
func (*Job) Payload ¶
Payload returns the raw wire envelope, which is a tag byte followed by the call body. It is opaque and must not be re-encoded.
func (*Job) Progress ¶
Progress reports how far along this attempt is, as a percentage.
Fire and forget: nothing answers it, and it never settles the job. It is a no-op when the scheduler did not advertise the side-channel capability, which is that capability's documented degradation.
The value is clamped to 0-100. The frame's range is part of the contract and the scheduler does not enforce it, so a value outside it would be stored and read back wrong rather than refused.
func (*Job) Publish ¶
Publish records a partial result for this attempt: a log line at level "result" whose message is empty and whose value rides in the frame's extra.
Logging and publishing share one frame, so a subscriber reading partials and an operator reading logs are reading the same stream.
func (*Job) RunKey ¶
RunKey is the id this durable run began under.
The job's own id for one that has only ever been retried in place — an ordinary retry, a requeue and a sleep wake all keep it. The id it started with for one an operator resurrected from the dead-letter queue, which is the single boundary where a job's id changes. Every idempotency key this attempt mints is built from it.
func (*Job) Sleep ¶
Sleep ends this attempt and reschedules the job to wake after d.
if err := job.Sleep(ctx, "cooloff", time.Hour); err != nil {
return nil, err
}
It returns ErrStepSlept once the deadline is committed, and the task body must propagate it: the row is written, the execution claim is released and the job is already pending at its deadline. Anything done past this point runs unclaimed and will run again on the wake.
On the attempt that wakes, the same call returns nil and execution carries on from there — which is what stops a job with three sleeps restarting the first one on the third wake.
The clock is read once, here. A replay keeps the deadline the row already holds rather than pushing it a full duration further out, so a job that crashes into its own sleep does not sleep forever.
type Option ¶
type Option func(*config)
Option configures a Worker.
func WithGRPCDialOptions ¶
func WithGRPCDialOptions(opts ...grpc.DialOption) Option
WithGRPCDialOptions appends raw dial options, for interceptors, a custom resolver, keepalive tuning, or anything else this package does not wrap. They are applied last and win over the options above.
func WithHandshakeTimeout ¶
WithHandshakeTimeout bounds the wait for a hello acknowledgement. A stream whose handshake does not complete inside it is torn down and retried.
func WithHeartbeatInterval ¶
WithHeartbeatInterval sets how often free capacity is reported. Heartbeats stop while the executor is draining.
func WithID ¶
WithID sets this executor's identity.
One live stream per id: a second attach under an id already attached is refused with ALREADY_EXISTS, and Worker.Run treats that as permanent rather than reconnecting into the same refusal. The default is unique per process, which is safe but tells an operator nothing; name your deployment instead.
func WithInsecureTransport ¶
func WithInsecureTransport() Option
WithInsecureTransport sends the token over an unencrypted connection.
For TLS, either flexiq-server terminates it (FLEXIQ_GRPC_TLS_CERT; pass WithTLS with its CA, and a client certificate for mTLS) or a proxy or a mesh in front of it does. This option is for the two hops that have no network to observe — a Unix-domain socket, and a loopback bind whose peers are on the same host — and for tests.
func WithLogger ¶
WithLogger sets where this package logs. Defaults to slog.Default.
It logs a rotation, a reconnect, a refused attach and a frame it does not recognise. All four are ordinary, and an executor that cannot say which one happened leaves an operator guessing.
func WithMaxMessageBytes ¶
WithMaxMessageBytes overrides the MaxMessageBytes cap in both directions.
Lowering it is the useful direction, for an executor that would rather fail early than buffer a large payload. Raising it past the server's own cap does not raise the server's.
func WithReconnectBackoff ¶
WithReconnectBackoff sets the reconnect schedule used after a transport failure. It doubles from min to max with jitter and resets on a completed handshake.
It does not apply to a stream the scheduler ended cleanly: that is a rotation, not a failure, and reconnecting is immediate.
func WithSDK ¶
WithSDK overrides the SDK name and version reported in the handshake, for a framework that wraps this package and wants its own name in the server's inventory.
func WithShutdownDrain ¶
WithShutdownDrain bounds how long a stream ending waits for the jobs it is already running.
Past it the connection closes anyway and whatever is still running is left to the scheduler's reaper — a handler that ignores its context must not be able to hang the process.
func WithSlots ¶
WithSlots sets how many jobs this executor runs at once. Defaults to GOMAXPROCS.
The scheduler reserves a slot before it writes a job frame, so it is designed never to oversend. A job that arrives with nothing free anyway is refused retryably rather than dropped.
func WithStepAckTimeout ¶
WithStepAckTimeout bounds how long a durable step waits for the scheduler to acknowledge its commit. Defaults to 30 seconds, the reference executor's own number.
The wait is bounded by the job's remaining time as well, whichever is shorter: waiting past the attempt's deadline only delays a reap the scheduler has already decided on. Running out is a **retryable** failure — an unconfirmed commit is indistinguishable from one that never happened, and the replay re-runs the step under the same downstream idempotency key.
func WithStepLimits ¶
func WithStepLimits(limits StepLimits) Option
WithStepLimits sets the caps a step commit is refused against before the round trip, for a scheduler configured away from the defaults.
The check that holds is the scheduler's either way. This one only buys an error that names the step and the value that failed, instead of one that arrives from the far side of the network.
func WithTLS ¶
WithTLS replaces the default TLS configuration, for a private CA or a pinned certificate. It clears a prior WithInsecureTransport, so the last transport option a caller passes is the one that holds.
func WithToken ¶
WithToken sets the bearer credential. Required, and it must carry the "execute" scope: a produce-scoped token cannot open an executor stream.
func WithTransportCredentials ¶
func WithTransportCredentials(creds credentials.TransportCredentials) Option
WithTransportCredentials sets the transport credentials directly, for a mesh or a credential type TLS does not cover. Like WithTLS, it clears a prior WithInsecureTransport.
func WithUserAgent ¶
WithUserAgent replaces the gRPC user agent, which by default names this client and its version.
type StepError ¶
type StepError struct {
// JobID is the attempt the step belongs to.
JobID string
// StepKey is the step's identity, "name#occurrence" or "name:key". Empty
// when the failure happened before one could be derived.
StepKey string
// Message is the reason, in whichever side's own words saw it.
Message string
// contains filtered or unexported fields
}
StepError is a step that could not be committed.
Test it with errors.Is against one of the verdicts above. A permanent one also answers to ErrFatal, so a task body that simply returns it settles the job the way the verdict says without the caller having to translate.
type StepLimits ¶
type StepLimits struct {
// MaxStepBytes is the largest encoded result one step may commit.
MaxStepBytes int
// MaxTotalBytes is the largest total across every committed step of a job.
MaxTotalBytes int
// MaxSteps is the most steps one job may commit.
MaxSteps int
}
StepLimits are the caps this client refuses a step commit against before the round trip.
A zero field takes its default; anything above the scheduler's hard ceiling is brought back to it, because a cap a caller can raise without bound is not a cap.
func DefaultStepLimits ¶
func DefaultStepLimits() StepLimits
DefaultStepLimits are the caps a worker uses when none are configured.
type Worker ¶
type Worker struct {
// contains filtered or unexported fields
}
Worker attaches to the executor door and runs the jobs it is given.
A Worker holds one gRPC connection. Build one per server, register the handlers, then call Worker.Run, which blocks.
func New ¶
New builds a worker for the server at target.
The target is a gRPC name: "host:port", or "unix:///run/flexiq.sock" for a Unix socket. No connection is made here — Worker.Run opens one — so New failing means the arguments were wrong, never that the server is down.
A token is required, and it must carry the "execute" scope.
func (*Worker) Close ¶
Close releases the connection. It does not stop a running Worker.Run; cancel its context for that.
func (*Worker) Handle ¶
Handle registers the handler for a task name.
Every registered name is advertised in the handshake, and nothing else is ever sent to this executor. Registration must happen before Worker.Run: the advertised list is fixed for the life of a stream, so a handler added later would be one the scheduler does not know exists.
func (*Worker) Run ¶
Run attaches and dispatches until the scheduler says stop or ctx ends.
It reconnects on its own. A clean stream end is a rotation — streams are bounded because a gRPC stream cannot be load balanced once started — and reconnecting from one is immediate and not an error. A transport failure backs off. A shutdown frame, a refusal that reconnecting would repeat, and a cancelled context each return.
Cancelling ctx begins a graceful drain: no new work is accepted, running jobs keep their contexts until the drain budget expires, and their results are still reported. Run then returns ctx.Err().