These commits are when the Protocol Buffers files have changed: (only the last 100 relevant commits are shown)
| Commit: | 109ac5c | |
|---|---|---|
| Author: | Julius Park | |
initial commit
| Commit: | 65c7d76 | |
|---|---|---|
| Author: | Alexander Belanger | |
Merge remote-tracking branch 'origin/main' into belanger/serverless-operator
| Commit: | 1925dc7 | |
|---|---|---|
| Author: | Julius Park | |
| Committer: | GitHub | |
feat: durable streams (#5072) * initial wip * fix tests, add default cursor behavior * rabbit publishes must be confirmed otherwise grpc error occurs, needs deslopping * forgot batch files * a variety of fixes/deslopping * some more fixes (first message race, tighten up sdk failures) * small generation fixes * fix lint * more linting and remove test that keeps failing * generate * refactors, etc * consolidate all pg reads into topic poller * retention stuff * more granular fixes * make partitions by tenantid so they can be dropped based on tenant retention * fix ordering bug * some validation changes etc * rabbit-ectomy, use direct to pg for grpc handlers * forgot some files * rework partitions etc * rework offset stuff * move migration * add feature flag * retention management * consolidate tests * clean up comments and fix load test * add changelog * add test files that were untracked * autovac settings * use default autovac for messages table * add reconnection * fix generate * fix v0config issue
| Commit: | 9678dd5 | |
|---|---|---|
| Author: | Alexander Belanger | |
Register serverless endpoints under their plain names. Each endpoint had a namespace UUID the operator prefixed onto every workflow, action and event it registered, carried in the healthcheck, trigger and first-frame bodies and confining the durable relay. The per-endpoint auth it was the hook for needs a larger refactor of namespace-scoped tokens on worker endpoints, so the prefixing goes: workflows and actions register exactly as a worker's do. Two endpoints of one tenant declaring the same action both serve it, like two workers: the routing cache keeps the enabled endpoints per action and picks one per delivery uniformly at random among those not known unhealthy, falling back to any enabled one, and a miss fails the task retryable without a repository lookup. The column, the proto fields (numbers reserved), the REST field, the miss query and the relay's confinement are removed, and the SDK REST clients follow. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK
| Commit: | 793dc97 | |
|---|---|---|
| Author: | Alexander Belanger | |
Name the contract timestamps by their unit. ServerlessHealthcheckRequest and ServerlessTriggerRequest carry the build time as Unix seconds; the field name now says so and the proto and header comments state the unit and the 5 minute freshness window, so an endpoint SDK cannot mistake the value for milliseconds. The field numbers are unchanged. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK
| Commit: | 2d799d9 | |
|---|---|---|
| Author: | Julius Park | |
merge
| Commit: | b4a480f | |
|---|---|---|
| Author: | Julius Park | |
| Committer: | Julius Park | |
initial wip
| Commit: | 8d1b4e3 | |
|---|---|---|
| Author: | Alexander Belanger | |
Merge remote-tracking branch 'origin/main' into belanger/serverless-operator Main added its own v1_0_157, so the serverless migration moves to v1_0_158. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK
| Commit: | 073f45f | |
|---|---|---|
| Author: | Gabe Ruttner | |
| Committer: | GitHub | |
fix(engine): look up satisfied durable entries by primary key (#5026) * fix(engine): look up satisfied durable entries by primary key The old join filtered branch and node after reading every entry for the task. * fix(engine): prune durable entry lookups to the task partition Pass each task's inserted_at into ListSatisfiedEntries and plan with those values so daily partitions that do not contain them are not locked. * fix(engine): bound durable entry lookups by task inserted_at range Replace the timestamp array with min/max bounds so the btree does not re-sort the array on every index rescan, and return to the plain primary-key join with materialized inputs instead of a lateral subquery. * simple feedback * drop plan stuff * fix(engine): pass task inserted_at on the operator action The DAG operator reads inserted_at from the assigned action instead of looking the task up. Keep the satisfied-entry external id typed as uuid so sqlc still emits uuid.UUID. * fix: signatures
| Commit: | 1ae390a | |
|---|---|---|
| Author: | Alexander Belanger | |
Add stream frames and a per-task streams flag to the serverless contract A serverless task that needs an engine stream (child run results, run events, durable events) cannot hold one from a fetch-based runtime, so the operator opens the stream on the task's behalf and relays it over the invocation websocket. The healthcheck catalog gains per-task options with a streams flag so the SDK can ask for a socket for a non-durable task, and the frame oneof gains stream_open, stream_message and stream_close. Existing frames and field numbers are unchanged; the TypeScript bindings are regenerated in the serverless SDK phase. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK
| Commit: | e990e1e | |
|---|---|---|
| Author: | abelanger5 | |
| Committer: | GitHub | |
feat(engine): engine-side gRPC operator support for out-of-process operators (#4910) * feat(engine): gRPC operator service for out-of-process operators Adds v1.OperatorService so an operator can run outside the engine and talk to it over gRPC with a tenant-scoped token: - Register (unary) upserts the operator by (tenant, name, kind = GRPC) and creates the worker for the connection, or resumes a previous worker of the same operator. Registration carries no workflows and no actions. - Listen (bidi) activates the worker for the lifetime of the stream: the first message is start{worker_id}, then heartbeats and action deltas flow in and AssignedAction messages flow out straight from the dispatcher fan-out. Deltas are applied incrementally (at most 1000 ids per message) and the scheduler is notified at most once per second while they arrive. - SendStepActionEvent and DurableTask delegate to the dispatcher; the durable stream is wrapped so the first register_worker message is checked for worker ownership before the dispatcher sees it. - Every RPC after Register, including Listen, carries hatchet-operator-id metadata validated against the token's tenant through a cached lookup. Repository: CreateWorkerOpts.OperatorId flows into the existing CreateWorker query and skips WORKER / WORKER_SLOT metering, since operator workers are excluded from those limit counts. AddWorkerActions / RemoveWorkerActions apply deltas in one transaction under the worker row lock and maintain the action hash as the XOR of per-action sha256 digests, so equal sets hash equal in any order and a delta costs O(chunk). The proto keeps importing the package-less dispatcher.proto for the shared action types (no SDK regeneration). SERVER_GRPC_OPERATORS_ENABLED (default false) gates service registration. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * feat(client): gRPC operator client with streamed action deltas Go client for v1.OperatorService in pkg/client, on the existing reconnecting stream machinery: Connect calls Register and opens Listen with the start message; heartbeats and action deltas go through retrySend so a dead stream reconnects and resumes the worker. AddActions / RemoveActions never block: a flusher coalesces pending ops (an add followed by a remove of an id that was never sent cancels out), chunks them to 1000 ids and sends every 250ms or as soon as a chunk is full. Flush waits for the queue to drain. The session keeps the desired action set and replays it as chunked deltas when a reconnect's Register did not resume the previous worker. PutWorkflow is a thin helper over the admin service that returns the derived action ids for the caller to AddActions. Harness-based e2e suite covering echo, durable memo, deactivation on close, bulk deltas with hash consistency, resume without replay, non-resume replay, and the operator id and worker ownership checks. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * fix(dispatcher): keep operator stream sends serialised and drains unblockable The shutdown drain ranges the session map and did a blocking send on each stream session's fin channel. A Listen handler that exited concurrently released its session and stopped selecting on fin, so a pointer the drain had already captured blocked it forever. Operator stream sessions now carry a done channel that Release closes, and the drain selects on fin or done through requestFin, so a released session is skipped. Sends on one stream are serialised by the per-worker send lock, but the lock was released by the caller's deferred Release as soon as its context ended, while the SendMsg goroutine was still blocked by flow control. The next caller then started a second SendMsg on the same stream, which gRPC forbids. The lock is now released by the send goroutine once SendMsg exits; a caller whose context ends returns without the lock, and the next caller fails fast with errFlowControlActive after the lock timeout instead of overlapping. AddOperatorStreamSession returns an OperatorStreamSession handle with Fin, Send and Release. Send shares the send serialisation with the action fan-out so an operator service can write its own protocol messages on the stream, and the optional wrap function converts each assigned action into the stream's server message type before it is encoded. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * fix(grpcoperator): acknowledge action deltas and retain them until committed The Listen protocol had no delta acknowledgement. A delta the stream had accepted but the engine had not committed (stream ended before the apply, database error, rejected chunk) was dropped by the client as soon as the send returned, and a resumed worker never learned about it: the client only replayed its desired set for a fresh worker, and repeating the same AddActions was a no-op against the desired set. Listen now returns a stream of OperatorListenResponse, a oneof of the dispatcher's AssignedAction and an OperatorActionsAck. Each delta carries a positive, strictly increasing sequence; the engine applies deltas in stream order and acks a sequence once the change is committed, and an ack for N confirms every lower sequence. A delta that cannot be applied, or whose ack cannot be written, ends the stream instead of leaving it unconfirmed. The dispatcher session wraps actions into the envelope and writes acks through the same per-stream send serialisation as the fan-out. The client keeps every sent chunk until its ack arrives and replays the unacked chunks in order on every reconnect. When the reconnect did not resume the worker it replaces the queue with a snapshot of the desired set, taken under the same lock AddActions and RemoveActions use, so a removal that races the replay is queued relative to the snapshot instead of being coalesced away with an add that was never sent. Replay runs as the reconnecting stream's replay callback, under sendMu before the stream is published, so no chunk is sent in between. A chunk that waits longer than the ack timeout hangs the stream up so the reconnect replays it. Flush now waits for the acks, so a nil Flush means the engine committed the set. Because Flush is called before Actions by the serverless link, the session runs its receive and heartbeat loops from connect until Close; the receive loop demultiplexes acks and actions, buffering actions in an inbox that the consumer Actions attaches drains. Starting the loops and attaching the consumer add to the WaitGroup under the mutex Close takes, which removes the Actions/Close race on the WaitGroup; Actions after Close is refused. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * fix(repository): hash worker action sets canonically and check the tenant on links The worker action hash was the XOR of per-action digests seeded from whatever the row held. A worker created with an initial set stored the sha256 of the ordered list while a worker built by deltas stored the XOR fold, so equal sets hashed differently across the two paths, and the XOR fold is linear, so a chosen set of ids could be solved to collide with any target hash. GetWorkerActionsByWorkerActionHash shares one worker's action list with every worker of the same hash, so both properties matter. hashActions is now the canonical digest: sha256 over the lower-cased, deduplicated ids sorted by byte order, each followed by ";". The delta transactions recompute the same digest in SQL from the linked rows (ComputeWorkerActionHash, sha256 over string_agg ordered with the C collation) while holding the worker's row lock, so the stored hash is always the digest of the links the transaction leaves behind. UpdateWorker, which only adds links, recomputes it the same way. At 100,000 linked actions one add chunk of 1,000 costs about 85 ms, a no-op re-add about 12 ms and a remove about 65 ms on the review database. Deltas no longer upsert the shared Action rows: existing rows are read, the missing ones are inserted in sorted order with ON CONFLICT DO NOTHING, and the resulting row ids are linked in sorted order. Concurrent replays of the same set by different workers therefore take no locks on existing Action rows, write no new row versions, and cannot deadlock. The row lock, the link and the unlink statements all carry the tenant: LockWorkerActionHash finds no row for a worker of another tenant, which is reported as an error wrapping pgx.ErrNoRows before anything is mutated, and the link and unlink statements join Worker and Action on the tenant so a cross-tenant pair affects zero rows. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * fix(grpcoperator): refuse a second register_worker on a DurableTask stream The durable stream wrapper checked worker ownership on the first register_worker only and passed every later message through, so a second register_worker naming another operator's worker in the same tenant reached the dispatcher, which rebinds the named worker's durable dispatcher without an operator check. The SDK's durable listener registers exactly once per stream and opens a new stream to register again, so a second register_worker is a protocol error: the wrapper now refuses it with InvalidArgument, whichever worker it names, and the stream stays bound to the worker it was authorized for. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * fix(grpcoperator): cap Listen streams and action links per operator Public gRPC admission had no cumulative limit: an authenticated tenant token could open any number of Listen streams for one operator, each a live worker exempt from the worker and slot metering, and stream deltas that grow the operator's action links without bound. The service now enforces two per-operator caps, refused with ResourceExhausted. WithMaxListenStreamsPerOperator bounds the Listen streams one operator holds open on this replica (default 100); a stream over the cap is refused before its worker is activated and the slot is returned when the stream ends. WithMaxActionsPerOperator bounds the action links held across all workers of the operator (default 1,000,000). The budget is read once per stream with CountOperatorWorkerActions and then tracked from the deltas the stream applies; the repository enforces it on the links a delta would newly create, inside the delta's transaction, so a delta over the budget links nothing and a delta that repeats actions the worker already holds (a replay after a resumed reconnect) never counts. Zero disables either cap. The caps are options on the service. Wiring them to server configuration (runtime.grpcOperatorMaxListenStreamsPerOperator and runtime.grpcOperatorMaxActionsPerOperator with env bindings, next to grpcOperatorsEnabled) is a change to pkg/config/server and cmd/hatchet-engine outside this branch's scope. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * fix(client): close operator clients and cancel durable sends Client gains Close, which closes the gRPC connection the client owns, so a caller that retires a client (the serverless grpclink when a tenant's token rotates or its last registration closes) releases the connection's goroutines instead of leaving them behind (correctness F14, performance F13). grpclink's closeClient calls it directly and the goleak regression test that had to skip while pkg/client exposed no Close now runs. DurableTaskListener.SendRequest takes a context and returns an error: the request queue is bounded, so a caller blocked on an unreachable engine is released when its context ends or the listener is stopped (Stop closes a per-generation stop signal that Start renews), and the listener's own request helpers drop their pending ack when the send is refused. SendMemoCompleted follows the same shape. grpclink's pump hands requests to the listener under a context closeAll cancels, so a stopped registration no longer holds the pump on a full queue (correctness F12). Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * fix(grpcoperator): expose admission caps in server config The per-operator Listen stream cap and action link budget were options on the operator service with no way to set them from configuration. Runtime now carries GRPCOperatorMaxListenStreamsPerOperator (default 100) and GRPCOperatorMaxActionsPerOperator (default 1,000,000), bound to SERVER_GRPC_OPERATOR_MAX_LISTEN_STREAMS_PER_OPERATOR and SERVER_GRPC_OPERATOR_MAX_ACTIONS_PER_OPERATOR next to SERVER_GRPC_OPERATORS_ENABLED, and both engine run paths pass them to grpcoperator.New through WithMaxListenStreamsPerOperator and WithMaxActionsPerOperator. Zero disables either cap (security F06). Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * fix(repository): recompute operator worker action hash from links UpdateOperatorWorkerActions stored the digest of the list it was given while only adding links, so a worker that already held links ended with a hash that named a subset of its linked set, and GetWorkerActionsByWorkerActionHash would then share the wrong action list with other workers of that hash. The call now locks the worker row with the tenant in the predicate (a worker of another tenant is refused and nothing is written), lower-cases and dedupes the ids, links through the tenant-checked LinkActionsToWorkerReturning, and recomputes the hash from the linked rows with ComputeWorkerActionHash inside the same transaction, as the worker repository's delta path does (correctness F17, security F02, F08). Tests check that the stored hash is the canonical digest of the final linked set after an additive call and after a call naming a subset, that it equals the delta path's hash for the same set, and that the tenant is checked. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * test(grpcoperator): synchronise the operator store fake fakeOperatorStore counted GetOperatorById calls and served its map with no synchronisation, while TestListenStreamCapIsPerOperator authorizes several Listen streams concurrently through it, so the race detector failed the package intermittently. The counter is atomic and the map is guarded, and the assertions read the counter through Load. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * refactor(grpcoperator): extract operatorsvc The gRPC operator service owns logic that is not gRPC's: upserting the operator row, creating or resuming its worker, activating the worker under a listener fence, applying action set deltas within the per-operator budget, and closing the session again. The in-process operator host needs all of it, and the serverless operator's engine link already re-implements it by hand, with the drift that comes from a copy: no resume, no budget, a separate dispatcher session key. internal/services/operatorsvc now holds that logic behind Service, Session and the two deliveries a session can have. A stream delivery encodes assigned actions onto a gRPC stream and counts against the per-operator stream cap; a handler delivery hands them to an operator running in this process and counts against nothing but the action budget, since it holds no stream. Both get one session id that is the dispatcher's session key and the worker row's listener fence, both notify the scheduler through the same throttle, and both close by pause, drain, deactivate. The gRPC Listen handler closes without the pause: its operator pauses itself before hanging up, and a stream that simply drops must leave the worker assignable for the next connection. internal/services/grpcoperator keeps what is genuinely the protocol: the metadata that names the operator, the start handshake, the message loop, the delta acks, and the durable stream it hands to the dispatcher unchanged. The dispatcher gains AddOperatorSession, the handler-backed twin of AddOperatorStreamSession, keyed on the session id its caller chose, and pkg/operator gains ActionHandler, the one method a session needs to deliver to. The test doubles for the stores and the dispatcher move to operatorsvc/operatorsvctest so the service's tests and the handlers' tests drive the same behaviour. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * feat(operatorsvc): durable invocations with ack-before-entry ordering An operator hosted in the engine drives a durable invocation over the dispatcher's channel-backed session rather than the DurableTask stream, and the handshake and ordering that goes with it is not obvious enough to leave to each host. Session.OpenDurable performs the register_worker handshake itself, so Recv only ever returns invocation traffic, and applies the ordering the operator expects: the engine delivers a completion for an already satisfied entry as soon as the invocation registers, which on a replay is before the operator has sent the wait_for that names it, so an entry is held until the ack carrying its ref has gone out. Responses that arrive during the handshake are held under the same rule, bounded at 256 so an engine that never acks cannot make the session buffer without end. As on the stream, one ack-bearing request is in flight at a time, because the engine keys its pending state by (task, invocation). This is a port of the serverless operator's engine link, which grew the ordering rule as a bug fix and will be deleted in favour of this once the in-process host lands. The gRPC DurableTask handler is unchanged: it still hands the raw stream to the dispatcher after the ownership check. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * feat(grpcoperator): PauseWorker An operator that wants to shut down cleanly needs the scheduler to stop assigning to its worker before it stops answering, and it needs to know the pause is committed before it starts draining. A message on the Listen stream would need an ack of its own and would not work while that stream is being torn down or reconnected, so pausing is its own unary call. OperatorService.PauseWorker takes {worker_id, paused}: the same call pauses and unpauses, which is what a worker that is resumed rather than replaced needs. It is authorised like SendStepActionEvent, by the operator metadata plus worker ownership. Register already clears the pause when it resumes a worker, so an operator that paused to drain and then crashed comes back assignable. The proto is unreleased, so this is an addition with no compatibility shim. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * feat(client): pause-then-drain on session close Closing an operator session used to hang the Listen stream up at once, which deactivates the worker while the actions the consumer is still working on are in flight; the engine only recovers those when the tasks time out. Close is now pause then drain: it pauses the worker through PauseWorker, so the scheduler stops assigning, waits for the task runs already handed to the consumer to be reported, and only then flushes the deltas and ends the stream. The session learns what is in flight from the contract it already has with the caller: a task run is outstanding from the moment it is handed to the Actions channel until the caller reports it completed, failed or cancelled. Cancels carry no work of their own and are never reported, so they are not tracked. The wait is bounded (30 seconds by default, WithDrainTimeout to change it) and WithoutDrain keeps the previous immediate hang-up, which is what a process that is going away anyway wants. Pause and Resume are also on the session, for a caller that drains on its own schedule rather than at close. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * fix(client): count an action in flight before handing it over Two ways the drain on Close could wait for reports that were never coming. An action was recorded as in flight after the send on the Actions channel returned, so a consumer that reported the outcome the instant it took the action cleared an entry that had not been added yet; the entry added afterwards was never removed and every later Close waited out the whole drain timeout. The action is now recorded before the hand-over, and taken back when the hand-over does not happen. The other way is a consumer that walks away: it cancels the context it passed to Actions, and the runs it was already holding will never be reported. Ending delivery now forgets them, so Close hangs up instead of waiting; the engine retries those tasks when they time out, which is what happened to them before this branch too. Also: report the drain timeout only when something really is outstanding, keep the operator name on the Listen handler's log lines, and say plainly that the notifier may publish one notification after stop. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * feat(operator): host contract Add the hosting contract to pkg/operator: Identity (tenant plus an existing operator row or a name and kind to upsert), OpenOpts, Host, Session, Registration, DurableChannel and the sentinel errors. It is what an operator needs from the engine whichever process it runs in; the engine-internal types (Operator, TaskEventWriter, SharedOperator) stay in operator.go and are documented as such. operatorsvc grows what the in-process host needs to implement it: - RegisterOpts.OperatorId registers an already claimed row (name and kind from the row, no upsert) and points the row's worker_id at the new worker, which is how ClaimOperators keeps seeing the assignment. - Session.SendStepActionEvent reports on the tenant-forged context path, filling the session's worker id and refusing another's. - The durable channel takes contexts on Send and Recv, uses the contract's error sentinels, stamps and passes worker_status through (the operator reports what it is blocked on; the worker is the session's), and reads the engine's responses on a pump goroutine so a single-threaded operator that sends while a response is undelivered no longer deadlocks against the engine's sequential delivery. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * feat(operator): in-process host over operatorsvc internal/operator/hostinproc implements pkg/operator.Host inside the dispatcher process. Open registers the operator through operatorsvc.Register (an existing row by id for claimed operators, an upsert by name otherwise), opens a handler-backed engine session whose id is both the dispatcher key and the worker row's listener fence, links the initial action set in engine-sized chunks and adds the worker to the host's heartbeat set. Deltas are applied synchronously, so Flush returns at once; events go through the session's tenant-forged path; OpenDurable is the engine session's handshake and hold buffer; PutWorkflow goes through the admin service and reads the derived action ids from the stored steps. One host-wide ticker writes a bulk UpdateWorkerHeartbeats for every open session, including those draining between Pause and Close. pkg/operator/operatortest carries an echo operator written against the contract only; it is opened here in a unit test and, in a later commit, over gRPC in the operator end-to-end suite. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * feat(operator): gRPC host pkg/operator/hostgrpc implements pkg/operator.Host over pkg/client's OperatorSession, a port of the serverless branch's grpclink onto the contract. It keeps one engine client per tenant (rebuilt when the token rotates, closed when replaced, released or on the host's Close), resolves tokens through a TokenSource (the TenantTokenExchange renamed; LocalExchange keeps the YAML file with referenced-file mtime tracking, StaticExchange the single-token case), refuses a session the engine authenticated as another tenant, and retries once after an Unauthenticated connect. A session runs the deliver loop that hands the client's assigned actions to the handler; a refused start is reported as a retryable failure because there is no requeue over gRPC. Pause is the unary PauseWorker RPC and Close is the client's pause-then-drain. The durable hub (task-indexed, bounded outbound queue, one ack-bearing request in flight, evictions and non-determinism errors reconstructed as responses) is mapped onto DurableChannel with context-aware Send and Recv; worker_status is dropped because the listener reports the awaited entries itself. Identity by operator id, a kind other than GRPC, ResumeWorkerId and WorkerName report ErrNotSupported: OperatorService has no such requests yet. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * refactor(operator): manager becomes a claimer over the in-process host Operator is now the contract's operator: an ActionHandler plus Start(ctx, Session) and Drain(ctx). SharedOperator no longer holds a repository or a worker id; it takes the session from Start and reports events, changes the action set (as a diff into AddActions and RemoveActions) and reads its worker id through it. TaskEventWriter keeps only what exists inside the engine: CancelTaskEvent, RegisterDurableTask, TriggerDAGStep and CancelDAGChildren. The DAG operator keeps its repository and writer and gains the lifecycle: Start registers the tenant's DAG actions on the session and starts the poller, Drain interrupts the runs in flight and waits for them. It is hosted in process only. pkg/operator/manager is replaced by internal/operator/claimer, which polls ClaimOperators as before but opens every claimed row through the in-process host (Identity{OperatorId}) and starts the operator on the session; a row that leaves the claim result is paused, drained and closed, and so is every row at Stop. Worker rows, heartbeats and the dispatcher's routing table are the host's now. It moved under internal/ because it depends on the in-process host and on operatorsvc, both engine-internal, and is constructed only by the engine. The dispatcher stops constructing the manager: listenForOperators, the om field and WithDAGOperatorDefaultSlots are gone, and with them the dispatcher -> manager -> operator import edge. run.go wires the claimer next to the dispatcher in both service sets and stops it before the dispatcher drains, while events can still be reported. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * fix(dagoperator): durable handshake through the session The DAG operator opened its durable invocations through RegisterDurableTask, sent register_worker itself and treated the first response as the ack without checking its type. The engine routes responses to the task as soon as it is registered, before the ack, so a completion restored for a resumed invocation could be taken for the ack and lost (F21, the race the serverless enginelink fixed for itself). The run now opens the invocation through Session.OpenDurable: the host does the handshake and holds any response that arrives before the ack, and delivers entries only behind the ack that names them. The dag drives a DurableChannel instead of the raw channel pair; the channel reads the engine's responses on its own, so the dag's send no longer has to drain responses to avoid deadlocking against the engine's sequential delivery, and awaitResponse's blocked-status report runs off a bounded Recv. The dag keeps its own buffer for an entry that races its ack, which the tests still exercise through an unordered test channel. RegisterDurableTask leaves TaskEventWriter and SharedOperator; the dispatcher still implements it for operatorsvc. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * test(operatore2e): host the echo operator over gRPC The echo operator in pkg/operator/operatortest is written against the contract only; internal/operator/hostinproc opens it in a unit test. This scenario opens the same handler through pkg/operator/hostgrpc against the harness engine: identity from the token's tenant, a workflow put through the session, a run driven to completion, and the host's teardown order (pause, drain, close) observed on the worker row. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * chore(operator): comments describe the hosts as they are The operatorsvc package doc names the in-process host that now calls it, ClaimOperators' doc no longer points at CreateOperatorWorker (the claimer registers workers through the operator service), and the DAG step index conversion says why its gosec finding is suppressed. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * docs(operator.proto): say why the heartbeat rides the Listen stream SDK workers heartbeat out of stream because the runtimes they target can block their event loop while a task runs and would starve an in-stream heartbeat. Operators are written against the Go SDK, which is not subject to that, so the heartbeat stays on the stream it keeps alive. The reasoning now lives on the OperatorHeartbeat message so it travels with the protocol. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * refactor(operator): the engine-internal writer names its tenant TaskEventWriter.CancelTaskEvent becomes CancelTaskEventCustom, named for what it is: the writer behind SendCancelledWithMessage, which reports a cancellation with a custom reason, distinct from the CANCELLED step action event every host offers and not on the gRPC surface. Every writer call now takes the tenant as an argument. The dispatcher's implementation no longer reads the tenant the gRPC auth middleware puts on a request context, so pkg/operator stops forging that context value with a copy of the middleware's key; the DAG step calls already carried the tenant and only pretended to need the context. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * refactor(operator): document the lastActions invariant, drop assertions UpdateWorkerActions keeps diffing against the set it last advertised: the engine would accept the whole set on every poll, but each send is a write, and the removes cannot be derived without the previous set. The map needs no lock because the operator calls the method from one goroutine at a time (its Start, then the poller Start launches once that call returned); the doc now says so instead of leaving it implied. The compile-time interface assertions go: the DAG factory, the hosts' Open and the test operator already return these types as the interfaces, so the compiler checks the same thing where it matters, and the repo does not use the idiom elsewhere. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * refactor(hostinproc): take the admin service and the workflow repository directly The in-process host no longer needs a wrapper that forwards to the admin service and the workflow repository: it asks for each on its own, an AdminService for the put and a WorkflowStore for the stored steps the action ids are read from, both satisfied by what the engine already builds. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * refactor(operator): constructors take variadic options for every argument operatorsvc.New, hostinproc.New and claimer.New follow the repo's New(opts ...Opt) shape: required collaborators arrive through With... options like the optional ones and are validated in New with an error that names the option to use, the way grpcoperator.New and dispatcher.New do. The Deps structs go away; the claimer registers factories one kind at a time through WithFactory. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * feat(repository): apply an operator actions delta in one transaction A delta carries adds and removes under one sequence number, and the ack for that sequence tells the operator the delta is committed, so its two sides must commit together. ApplyWorkerActionsDelta replaces AddWorkerActions, AddWorkerActionsWithinBudget and RemoveWorkerActions: one row lock, adds then removes, the budget check before any remove is applied so a refused delta rolls back whole, and the action hash recomputed once from the resulting set. operatorsvc applies a delta with one call and charges the budget with both counts. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * fix(repository): meter operator workers like SDK workers A gRPC operator's workers are the tenant's workers: CreateNewWorker meters them against the WORKER and WORKER_SLOT limits, and the limit queries count them. The one exception is the DAG operator, whose workers the engine creates for every tenant with DAG workflows as infrastructure; those stay out of both, keyed on the operator's kind (a NOT EXISTS against v1_operator) rather than on operatorId, so what CreateNewWorker meters is exactly what the queries count. CreateWorkerOpts names the operator's kind next to its id for that decision. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * refactor(client): operator client and stream helpers get packages of their own pkg/client/streaming holds the reconnecting stream, the listen loop and the error classifier every listener in pkg/client shares, exported. On top of it, pkg/client/operatorclient is the OperatorService client: Connect, the session, the action delta queue and the durable stream opener, with a New(conn, opts...) constructor that takes the token and metadata it needs. It does not import the legacy package, so client.Operator() stays as a thin accessor built over the client's own connection. hostgrpc now speaks operatorclient. The legacy package is named in one file, for the two things only it has (the dial from a token and the environment, and the durable task listener the Go SDK worker shares), and the host builds and owns the listener since the client session no longer does. hostgrpc.New takes variadic options like the other constructors: WithTokenSource is required, WithLogger and WithClientFactory are not; the local exchange's logger option is WithExchangeLogger. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * docs(self-hosting): document the operator settings SERVER_GRPC_OPERATOR_MAX_LISTEN_STREAMS_PER_OPERATOR, SERVER_GRPC_OPERATOR_MAX_ACTIONS_PER_OPERATOR and SERVER_DAG_OPERATOR_DEFAULT_SLOTS join SERVER_GRPC_OPERATORS_ENABLED in the configuration table. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * feat(operator): pause is a message on the Listen stream The Listen stream is stateful, so the pause belongs on it rather than in a unary RPC. The client sends OperatorListenRequest.pause and waits for the OperatorPauseAck the server answers with once the pause is committed; from that ack on nothing is delivered on the stream. The server commits the pause in this order: the dispatcher session stops delivering first, then the row is written, then the ack goes out, so an action the scheduler assigned before it observed the pause is returned to the queue the way a failed send is (subscribedWorker.paused, honoured for starts and not for cancels, which a draining operator may still need). Resuming lifts the pause in the opposite order and is acknowledged the same way. The pause is stream state: Register clears it when it resumes a worker, so a paused client sends it again on every new stream, right after start and before its deltas. The unary PauseWorker RPC and its messages are gone; pkg/client/operatorclient keeps Pause, Resume and the pause-then-drain Close over the stream, and the in-process host pauses through the same operatorsvc.Session.Pause, so both delivery kinds stop the same way. The e2e pause scenario now asserts that nothing reaches the operator after the pause ack and that the run held back runs once after the resume. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * feat(operator): the contract names entries acked outside it and a session's end Two additions to pkg/operator, each for a review finding fixed on top of it. DurableChannel.ExpectEntry(branchId, nodeId) registers an entry whose completion the operator awaits but which no request on the channel acknowledged. The engine-internal writer creates DAG children directly and returns their refs, so no trigger_runs ack ever names them; the channel's ack-before-entry rule would otherwise hold their completions forever. In process it marks the ref as acknowledged (a held completion is queued at once); over gRPC it awaits the entry on the listener, once per entry. Session.Done and Session.Err surface the end of a session to the caller that owns it: Done closes after Close or when the host gave up on the session, and Err says which. Every implementation closes Done on Close today; the gRPC host's supervision of terminal failures builds on it. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * fix(safeclient): speak HTTP/1.1 only so the idle pool bounds every origin Performance R7. HTTP/2 connections live in net/http's h2 pool, which MaxIdleConns does not bound, so a fleet polling many distinct origins retained one connection per origin for as long as the origin was revisited: 500 origins left 500 connections open. The Sender now speaks HTTP/1.1 only, with ForceAttemptHTTP2 off, an empty TLSNextProto map and ALPN offering http/1.1 alone, so the custom dialer safeurl installs can never see h2 frames. Every retained connection is then in the transport's idle pool, bounded by the new Config.MaxIdleConns, MaxIdleConnsPerHost and IdleConnTimeout (256, 4, 90s by default, the names the serverless branch already uses), and Sender.CloseIdleConnections reaches the transport safeurl actually installed. Regression: TestDeliver_IdlePoolBoundsRetainedConnectionsAcrossOrigins (500 origins, MaxIdleConns 8: retained 500 over h2 before, 8 after) and TestDeliver_NegotiatesHTTP1AgainstHTTP2Origin. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * fix(dagoperator): expect directly triggered children and register one condition at a time Correctness R09 and R10, found by running the interpreter over the engine's own durable channel instead of the raw test channel. R09: a child the interpreter creates through the direct TriggerDAGStep callback has no trigger_runs request on the channel, so no ack ever names its entry and the channel held the child's completion for good: the run never advanced past its first asynchronous child, and the blocked-status recovery resent a completion the channel held again. The interpreter now registers every triggered, not yet satisfied step (and the on-failure step) with DurableChannel.ExpectEntry as soon as the callback returns its ref, which stands in for the ack; a completion that arrived first is released at once. R10: a task with a skip and a cancel condition, or several ready tasks with conditions, sent their wait_for registrations back to back, and the channel admits one ack-bearing request until its ack arrives, so the second was refused with ErrRequestInFlight and the run failed. The emitter now registers one condition per pass: a task whose registration must wait for a pending ack stays pending, the main loop consumes the ack (which is what assigns the registration its node id, in FIFO order), and the next pass registers the next one. Log order is unchanged; the trigger gate on pending acks already existed. The racing-ack test in dag_test.go is adapted: with one registration in flight the skip fires on its own ack, so the wait watch is never sent. Regression tests, over operatorsvc's channel and the dispatcher double: TestDagOverEngineChannelCompletesDirectlyTriggeredChild (failed with context deadline exceeded) and TestDagOverEngineChannelRegistersConditionsOneAtATime (failed with ErrRequestInFlight). Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * fix(operatorsvc): bound what the durable pump retains and deliver a completion once Correctness R07. After the register-worker ack the in-process channel's pump queued every engine response for Recv, and held every completion whose ref no ack had named, with no bound: 10,000 completions were accepted with no Recv, inside the engine process. RetainedResponseLimit (four handshake holds, 1024) now bounds queued plus held responses for the life of the invocation. Past it the channel fails: what it retained is dropped, Send and Recv return ErrChannelClosed, and the pump keeps discarding until the operator closes the channel, so the engine's delivery goroutine is never left blocked on a failed channel. The DAG operator reports such a run as a retryable failure. A completion is delivered once. The engine resends satisfied entries the operator reports as awaited, and a repeat used to be held forever under the ack-before-entry rule; delivered refs are remembered and a repeat is dropped, which also keeps expectations for delivered entries inert, as the contract says. The gRPC channel's ExpectEntry awaits an entry once however often it is registered, since a second listener callback would strand the first. Regression tests: TestOpenDurableRetainedResponsesAreBounded (Recv timed out against the old channel), TestOpenDurableDeliversACompletionOnce, TestOpenDurableExpectEntryStandsInForTheAck and, over gRPC, TestDurableChannelExpectEntryAwaitsTheEntry. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * fix(repository): one action budget per operator, an unambiguous hash refreshed per sequence, a fenced pause Four findings on the worker action set, all in the delta transaction and the session around it. Correctness R02 and R11, security R01: the per-operator action cap was a per-session snapshot tracked from that session's own deltas, read and written without synchronization from concurrent AddActions (a Go data race under the in-process host) and never shared: N sessions of one operator could each consume the whole cap. The budget is now one transactional check. "Worker"."actionCount" (new column, an invariant every link-mutating path maintains: the delta settles it by added minus removed under the worker row lock, CreateNewWorker sets the initial size, UpdateWorker recounts) makes the operator's total a sum over its workers, O(workers) never O(links). ApplyWorkerActionsDelta takes the operator-wide cap, locks the v1_operator row after the worker row (always in that order), sums the counts and rolls the whole delta back past the cap with *ActionBudgetError carrying the real totals. The session keeps no budget state at all; the fake store mirrors the operator-wide check. Correctness R03: the pause a closing session wrote, and Session.Pause itself, were not fenced on the listener session id, so closing a superseded session paused its successor. SetWorkerPausedForListener writes only while the session is the one on the row; a superseded write is a debug line, the newer session owns the worker's scheduling state. The explicit PauseWorker RPC path stays unfenced. Security R02: the canonical hash joined ids with ";" while an id may contain ";", so two different sets could hash equal and the scheduler, which treats equal hashes as equal sets, could give a worker another set's actions. Go and SQL now hash each sorted id as a 4-byte big-endian length prefix followed by its bytes. The migration recomputes the hash of active workers in place and clears it on inactive ones, which refresh when a session resumes them (OpenSession refreshes a NULL hash before activation, so a worker whose previous session died between a delta and its refresh comes back correct). Performance R2: every delta chunk re-sorted and rehashed all of the worker's links, 500 million row visits for a million-action replay, which timed out at 265,000. The delta now clears the hash and the session refreshes it once per delta sequence: in the throttled notifier before the scheduler notification that ends the window (only when a delta changed the links, on a detached bounded context), and on Close when a refresh is still owed. The scheduler already reads a worker without a hash through the join, which now also covers a pending refresh. Measured on the reviewer's dense fixture: 100k replay 10.3s to 3.1s (one refresh 0.23s); 1M replay from a 30s timeout at 265k to 26.0s (one refresh 2.3s cold, the link insert itself is 15.7s of it); one removal on a 1M worker from 715ms to 1ms plus a 1.0s refresh. The single 1M refresh is still over a second, so the alternative is recorded here: a (workerId, actionSetVersion) cache key, bumped by the settle statement, removes the full-set read entirely at the price of one cache entry per worker instead of per identical set. UpdateOperatorWorkerActions had no caller and is deleted with its test. The e2e hash reads wait for the end-of-sequence refresh. Regression tests: TestSessionActionBudgetSpansSessions, TestConcurrentDeltasThroughHost (race), TestOperatorActionBudgetSpansWorkers, TestOperatorActionBudgetUnderConcurrentDeltas; TestSessionSupersededCloseDoesNotPauseSuccessor, TestSupersededCloseLeavesSuccessorAssignable, TestPauseWorkerForListenerIsFenced; TestHashActionsIsCanonical, TestWorkerActionHashEncodingIsUnambiguous; TestWorkerActionHashRefreshFollowsDeltas, TestSessionRefreshesActionHashPerWindow, TestOpenSessionRefreshesPendingHash. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * fix(operatorclient): a chunk staged before a replay is never sent after it Correctness R01. The flusher stages a chunk (moves it from pending to unacked) before it takes the stream's send lock. A reconnect that owned the lock meanwhile replayed the stream: a fresh worker got a snapshot of the desired set, then the stale staged chunk went out behind it and undid part of the snapshot (wire order [2 1]; the worker lost an action while Flush reported success), and a resumed worker got the chunk twice, once from the replay and once from the flusher. The queue now carries a stream generation, bumped under its mutex by every replay. A chunk is staged with the generation it saw, and the flusher's send, which runs under the same send lock as the replay, checks the generation first: if a replay ran in between, the chunk is not sent and counts as done, because the replay either resent it from unacked (resumed) or dropped it for the snapshot (fresh). Regression tests: TestOperatorSessionFreshReplayDropsStagedRemoval and TestOperatorSessionFreshReplayDropsStagedAdd (the server set diverged from the desired set and the wire carried [2 1]) and TestOperatorSessionResumedReplaySendsStagedChunkOnce. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * fix(hostgrpc): supervise the client session and keep tenant clients referenced Four findings on the out-of-process host. Correctness R05: a terminal failure of the client session (a revoked token, an Unauthenticated Register on reconnect) ended the client's receive loop and the host only logged it, leaving the tenant registered and unserved for good; the token source was never asked again. The host session now supervises its client session: when the client's action stream ends while the session is not closing, it closes the durable hub (open channels fail and their invocations retry), closes the old client session and reconnects through the host with jittered backoff (1s to 30s), which asks the source for the tenant's current token, verifies the tenant, restores the action set the session advertised and resumes delivery under the new registration. It gives up only on a failure no retry fixes (ErrNoToken, a tenant mismatch, ErrNotSupported), reported through the contract's Done and Err. Correctness R04: ReleaseTenant, token rotation and the host's Close closed the tenant's cached client unconditionally, so an old registration's teardown closed the connection a newer registration of the same tenant had just opened over. Cached clients are reference counted: a session holds one reference from connect to Close (or to giving up); eviction only stops the client being handed out, and an evicted client is closed by its last release. Host.Close still closes everything. Correctness R06: the session reported outcomes under the worker its first Open registered, even after a fresh reconnect gave the client a new worker because the engine no longer had the old one, so the engine refused the reports. Registration follows the client's current worker; the worker current at each delivered start is recorded per task run, and a report is stamped with it while it is still the current worker and with the current worker otherwise, since the engine registers a fresh worker exactly when it cannot resume the previous one. A report naming a worker the session used to be is rewritten the same way; a foreign id passes through for the engine to refuse. Performance R8: delivery called the handler serially straight off the client's action channel, so a start that blocked on the operator's flow control held every later message, cancels included, while the client's inbox grew without bound. Delivery now reads the channel eagerly: cancels go to the handler on their own goroutine at once, starts go through a bounded queue (twice the worker's slot units, at least 64) served by one goroutine, and a start that finds the queue full is refused with a retryable failure report instead of being held. Regression tests: TestSessionReopensAfterTerminalFailure, TestSessionGivesUpWithoutToken, TestSessionCloseDuringBackoff and, over a real listener, TestLoopbackSessionRecoversWithRotatedToken; TestReleaseTenantWaitsForOpenSessions; TestSessionReportsUnderCurrentWorker and TestLoopbackFreshReconnectReportsUnderNewWorker; TestCancelsBypassBlockedStarts. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * fix(repository): keep the action hash encoding and drop the hash migration The worker action hash keeps the per-id framing main has always used: sha256 over the sorted, deduplicated, lower-cased action ids, each followed by ";". hashActions in Go and ComputeWorkerActionHash in SQL compute it byte for byte the same way, so a worker created with an initial set and a worker built by deltas hash equal for the same set. The length-prefixed encoding is dropped, and with it the two "actionHash" rewrites in migration v1_0_154, which now only adds "actionCount" and backfills it for operator workers. A rolling deploy needs no hash step: a worker registered by an older engine keeps a hash the new engine computes the same way, and any hash the scheduler does not recognise is a group of its own, so distinct hashes never misroute. The same-tenant ambiguity of a ";" inside an id is closed at validation instead, in the commit that follows. The unambiguous-encoding test becomes a Go/SQL parity test over a fixed set of action sets, and the unit test pins the digest of a sorted set to sha256 of "svc:a;svc:b;". Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * fix(validator): reject semicolons in action ids ParseActionID refuses a ";" anywhere in an action id. The worker action hash frames each id with a semicolon, so an id carrying one could make two different action sets of the same tenant hash equal, and the scheduler treats equal hashes as equal sets. ParseActionID is the one function every registration path goes through: the actionId validator tag (worker registration and update, operator action deltas, workflow step actions and the worker list filter) calls it, and the v0 and v1 PutWorkflow handlers call it directly on each step action before the step is stored. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * refactor(repository): rename "actionCount" to "operatorActionCount" The column added by migration v1_0_154 counts the "_ActionToWorker" rows a worker holds so the per-operator action budget can be summed over the operator's workers instead of counting their links. Only operator workers are read and only they are backfilled; an SDK worker's count is zero until it registers and nothing reads it. The name says so. Semantics are unchanged: every path that links or unlinks actions maintains it, and the budget check sums it under the operator's lock. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * chore(generate): format the configuration table and align the operator protoc stamp The configuration-options table added on this branch was not prettier formatted, so the docs generator reflowed its column widths and both the generate diff check and the docs lint job failed on it. The table is now written the way prettier writes it. operator_grpc.pb.go carried a protoc v5.29.3 stamp from a local protoc; CI generates with v5.29.6, like every other generated file here. The stamp line is aligned by hand rather than regenerated with the local toolchain. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * fix(operatorsvc): admit orchestrator action ids on DAG operator sessions The DAG operator registers the orchestrator action ids of the tenant's DAG workflows, "<workflow>_orchestrator". They are not "service:verb" action ids on purpose: without a ":" no SDK or GRPC worker can register one. On main the operator wrote them through a dedicated repository method that never ran the actionId validator. On this branch the in-process host routes them through Session.ApplyDelta, which validated every delta with the actionId tag, so every refresh failed with InvalidArgument from the first one at engine start, the DAG operator's worker never held an action, and no multi-task workflow of a tenant with dag_operator enabled ever started. That is the failing python test (3.14, true, true) cell: test_cancellation, and the durable spawn-DAG tests, hang until the job times out. Delta validation is now scoped by the operator's kind. A DAG operator session admits orchestrator ids and nothing else; every other kind admits action ids and nothing else. A DAG operator is hosted in process only and its row is never upserted through a registration, so the kind is the engine's own claim and a remote operator cannot hold an orchestrator action. The id form lives in one place, repository.DAGOrchestratorActionId, which the workflow store builds with and the service checks against; the form has no ";" either. The hostinproc tests that seeded DAG rows with "dag:a" style ids had hidden the mismatch; they use orchestrator ids now, and an operatorsvc test reproduces the failure on both kinds. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * feat(repository): add operator leasing and worker limit exemption columns What an operator is and who keeps it alive are separate axes. v1_operator.kind stays what it is (DAG, the engine-internal DAG operator; GRPC, a contract operator hostable in process or out of process) and a new v1_operator.leasing enum says who keeps the row alive: MANAGED rows are assigned to a dispatcher by ClaimOperators and built from a factory inside the engine, SELF rows keep themselves alive through a Listen stream out of process or their own leaser in process. SERVERLESS is retired: the serverless operator is GRPC with SELF in both modes, and the migration folds any SERVERLESS rows into GRPC (the enum value stays; databases migrated by the serverless branch carry it). Rows are unique per (tenant, name, kind) for every kind, replacing the two partial indexes that only covered the kinds upserted by name. Duplicates are renamed, not deleted, before the index is built. Limit exemption is a hosting fact decided at worker creation rather than an attribute of the operator row, so "Worker"."exemptFromLimits" is added, backfilled for the DAG operator's workers, which the limit queries left out before. The queries switch to the column in a following commit. The migration is v1_0_159: belanger/serverless-operator, which stacks on this branch, owns v1_0_155 to v1_0_158. It runs in goose's transaction: the enum type is created here rather than extended, and only a value added by ALTER TYPE ... ADD VALUE is unusable in the transaction that added it. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * refactor(repository): exempt workers from limits by flag Whether a worker counts against the tenant's WORKER and WORKER_SLOT limits used to be derived from the kind of the operator it backs (DAG workers were exempt). That tied metering to the operator row, which no longer says who hosts the operator. It is now an explicit CreateWorkerOpts.ExemptFromLimits, stored on the row, that only the in-process operator host sets: every worker it creates is engine infrastructure that runs whether or not the tenant runs workers of its own, whatever kind or leasing its row has. Wire-registered workers, over OperatorService or from an SDK, are metered. The limit queries in workers.sql and tenant_limits.sql read the flag off the worker instead of joining v1_operator on kind, so what CreateNewWorker meters is what they count. CreateWorkerOpts.OperatorKind had no other reader and is removed. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * feat(operator): separate operator kind from leasing Registration and claiming now run on the leasing axis rather than the kind list. ClaimOperators claims every MANAGED row, whatever its kind; the claimer resolves the factory by kind, except for GRPC rows, which it resolves by the row's name through WithNamedFactory, since the kind only says the row is a contract operator and the name says which one the engine should build. A claimed row with no factory is logged once and left alone. UpsertGRPCOperator is replaced by UpsertOperator(tenant, name, kind, leasing) over the plain unique index; a repeat registration takes the leasing it names, so a row the engine was leasing that registers itself leaves the claim set on the next poll, and the other way round. operatorsvc.RegisterOpts carries Leasing; a named registration must say who keeps the row alive. Wire registration (grpcoperator Register, hostgrpc) is GRPC and SELF, and AuthorizeOperator, behind the hatchet-operator-id metadata, requires kind GRPC and leasing SELF: an engine-leased row can never be driven over the wire. operator.Identity gains Leasing, defaulting to SELF; hostinproc passes it through and marks every worker it creates exempt from limits. The DAG operator's row is created MANAGED, and a race creating it is tolerated now that the name is unique. The unused CreateOperatorWorker query and repository method, which predate the claimer registering workers through the operator service, are removed. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * refactor(operator): name the leasing manager The v1_operator column that says who keeps an operator alive is now leasing_manager, of type v1_operator_leasing_manager with the values SELF and DISPATCHER. DISPATCHER replaces MANAGED: the dispatcher claims the row through ClaimOperators and builds the operator from a factory, a SELF row keeps itself alive. The Go side follows: sqlcv1.V1OperatorLeasingManager, the LeasingManager field on Identity, RegisterOpts, CreateOperatorOpts and UpsertOperatorOpts, and the comments that described managed rows. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * chore(migrate): one migration for the gRPC operator The three migrations this branch added (the GRPC enum value and its partial index, operatorActionCount, then leasing and exemptFromLimits) collapse into 202609111300000_v1_0_154.sql. The file runs with NO TRANSACTION so the GRPC value is committed before anything uses it; the partial GRPC index is never created, only the (tenant, name, kind) unique index; and the SERVERLESS rewrite is gone, since that value never exists on this branch. The goose id is fifteen digits: main's 202609101223115_v1_0_153.sql is, and goose refuses a version that sorts below one already applied. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * Renumber the operator migration to v1_0_155 with a fourteen-digit goose id. Main remapped v1_0_153 to a fourteen-digit id and added v1_0_154, so this file moves to 155 and drops the note about matching a fifteen-digit version. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * Address review comments on the operator migration, queries and naming. Migration: the "Worker" columns are added in one re-runnable statement at the top so the table lock is released before the backfills, the header states that as the reason for running without a transaction, and the comments are cut down. Queries: the action hash is recomputed and stored by one RefreshWorkerActionHash query, and the fenced pause goes through UpdateWorker with an optional listener session id instead of its own query. UpdateWorker now also scopes by tenant. Naming: isExemptFromLimits, is_paused on OperatorPause and OperatorPauseAck, CancelTaskWithReason, LinkNewActionsToWorker and UnlinkActionsFromWorker, and the durable goroutines are a receive loop (engine side) and a send loop (client hub) instead of a pump. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * Realign struct literals after the isExemptFromLimits rename. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * Constrain DAG operators to dispatcher leasing and make the down migration atomic. v1_operator gets a check that a DAG row is leased by the dispatcher. The down migration is sent as one query, so Postgres runs it as a single transaction and the drops need no IF EXISTS. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * Serve the operator service over ConnectRPC. The Listen and DurableTask handlers take connect's bidi streams and run over a transport-free half so their tests keep driving fake streams; the operator id is read from the request header through the handler call info; errors are connect errors with the same codes. The dispatcher's operator stream session takes the stream's guarded sender and wraps assigned actions itself, replacing the wrap hook, and the durable session is entered through the dispatcher's DurableTaskWithReceive so the ownership check on the first message stays in the operator service. The operator client stays on grpc-go, and the operator e2e suite passes against the connect server. Also fixes the two admin trigger paths that main left on the removed status package (#4755 landed beside #4987). Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK * Fix the load-tagged partition test for the two-value CreatePartitions. Main's #5002 changed the query to :one and left this test on the old signature. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK --------- Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com>
| Commit: | 376470c | |
|---|---|---|
| Author: | Alexander Belanger | |
Merge branch 'belanger/grpc-operator' into belanger/serverless-operator
| Commit: | d3ac62a | |
|---|---|---|
| Author: | Alexander Belanger | |
Address review comments on the operator migration, queries and naming. Migration: the "Worker" columns are added in one re-runnable statement at the top so the table lock is released before the backfills, the header states that as the reason for running without a transaction, and the comments are cut down. Queries: the action hash is recomputed and stored by one RefreshWorkerActionHash query, and the fenced pause goes through UpdateWorker with an optional listener session id instead of its own query. UpdateWorker now also scopes by tenant. Naming: isExemptFromLimits, is_paused on OperatorPause and OperatorPauseAck, CancelTaskWithReason, LinkNewActionsToWorker and UnlinkActionsFromWorker, and the durable goroutines are a receive loop (engine side) and a send loop (client hub) instead of a pump. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK
| Commit: | e640cb0 | |
|---|---|---|
| Author: | Alexander Belanger | |
Merge branch 'belanger/grpc-operator' Brings the ten #4910 review fixes and the seven round-2 fixes of the gRPC operator branch onto the serverless operator. Conflicts and adaptations: - pkg/operator/safeclient/safeclient.go: this branch's newTransport() and withDefaults() plumbing (shared by the SSRF and InsecureDestinations senders) with grpc-operator's HTTP/1.1-only transport settings (ForceAttemptHTTP2 off, empty TLSNextProto, ALPN http/1.1 alone, TLS 1.2 minimum) and its transport_test.go verbatim. - configuration-options.mdx: both branches' settings tables, the gRPC operator's row replacing this branch's shorter one. - cmd/hatchet-serverless-operator: hostgrpc.WithExchangeLogger for the local exchange, hostgrpc.New(WithTokenSource, WithLogger) returning an error. - Migration 20260911100000_v1_0_154.sql is renamed v1_0_157: goose orders it after this branch's 20260908110000..130000 (v1_0_154..156) by timestamp already, the rename keeps the version labels unique and in chain order. Its version number is unchanged. - CreateWorkerOpts.OperatorKind's oneof gains SERVERLESS, the kind the in-engine host registers under; grpc-operator's validator knew only HTTP_API, DAG and GRPC. - The serverless fakes implement Session.Done/Err and DurableChannel.ExpectEntry; the e2e process test builds the gRPC host through its options. - The runner supervises each registration's session: Done closing with an error means the host gave up for good (ErrNoToken, a tenant mismatch, ErrNotSupported; transient failures and token rotation the gRPC host recovers itself), so the registration is detached and closed and a new one opened while the tenant still owns units, with maintenance retrying an open that fails. Regression: TestRegistrationReopensAfterHostGivesUp. Metering: operator workers are now metered like SDK workers with only kind DAG exempt, so SERVERLESS workers count against WORKER and WORKER_SLOT in both hosting modes. Kept as merged; to be revisited. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK
| Commit: | 7e84168 | |
|---|---|---|
| Author: | Alexander Belanger | |
feat(operator): pause is a message on the Listen stream The Listen stream is stateful, so the pause belongs on it rather than in a unary RPC. The client sends OperatorListenRequest.pause and waits for the OperatorPauseAck the server answers with once the pause is committed; from that ack on nothing is delivered on the stream. The server commits the pause in this order: the dispatcher session stops delivering first, then the row is written, then the ack goes out, so an action the scheduler assigned before it observed the pause is returned to the queue the way a failed send is (subscribedWorker.paused, honoured for starts and not for cancels, which a draining operator may still need). Resuming lifts the pause in the opposite order and is acknowledged the same way. The pause is stream state: Register clears it when it resumes a worker, so a paused client sends it again on every new stream, right after start and before its deltas. The unary PauseWorker RPC and its messages are gone; pkg/client/operatorclient keeps Pause, Resume and the pause-then-drain Close over the stream, and the in-process host pauses through the same operatorsvc.Session.Pause, so both delivery kinds stop the same way. The e2e pause scenario now asserts that nothing reaches the operator after the pause ack and that the run held back runs once after the resume. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK
| Commit: | 579bb34 | |
|---|---|---|
| Author: | Alexander Belanger | |
docs(operator.proto): say why the heartbeat rides the Listen stream SDK workers heartbeat out of stream because the runtimes they target can block their event loop while a task runs and would starve an in-stream heartbeat. Operators are written against the Go SDK, which is not subject to that, so the heartbeat stays on the stream it keeps alive. The reasoning now lives on the OperatorHeartbeat message so it travels with the protocol. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK
| Commit: | 73267fa | |
|---|---|---|
| Author: | Alexander Belanger | |
Merge branch 'belanger/grpc-operator' into belanger/serverless-operator # Conflicts: # internal/services/dispatcher/operator_sessions.go # internal/services/dispatcher/operator_sessions_test.go
| Commit: | ed37b94 | |
|---|---|---|
| Author: | Alexander Belanger | |
feat(grpcoperator): PauseWorker An operator that wants to shut down cleanly needs the scheduler to stop assigning to its worker before it stops answering, and it needs to know the pause is committed before it starts draining. A message on the Listen stream would need an ack of its own and would not work while that stream is being torn down or reconnected, so pausing is its own unary call. OperatorService.PauseWorker takes {worker_id, paused}: the same call pauses and unpauses, which is what a worker that is resumed rather than replaced needs. It is authorised like SendStepActionEvent, by the operator metadata plus worker ownership. Register already clears the pause when it resumes a worker, so an operator that paused to drain and then crashed comes back assignable. The proto is unreleased, so this is an addition with no compatibility shim. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK
| Commit: | e97d912 | |
|---|---|---|
| Author: | Alexander Belanger | |
| Committer: | Alexander Belanger | |
fix(grpcoperator): acknowledge action deltas and retain them until committed The Listen protocol had no delta acknowledgement. A delta the stream had accepted but the engine had not committed (stream ended before the apply, database error, rejected chunk) was dropped by the client as soon as the send returned, and a resumed worker never learned about it: the client only replayed its desired set for a fresh worker, and repeating the same AddActions was a no-op against the desired set. Listen now returns a stream of OperatorListenResponse, a oneof of the dispatcher's AssignedAction and an OperatorActionsAck. Each delta carries a positive, strictly increasing sequence; the engine applies deltas in stream order and acks a sequence once the change is committed, and an ack for N confirms every lower sequence. A delta that cannot be applied, or whose ack cannot be written, ends the stream instead of leaving it unconfirmed. The dispatcher session wraps actions into the envelope and writes acks through the same per-stream send serialisation as the fan-out. The client keeps every sent chunk until its ack arrives and replays the unacked chunks in order on every reconnect. When the reconnect did not resume the worker it replaces the queue with a snapshot of the desired set, taken under the same lock AddActions and RemoveActions use, so a removal that races the replay is queued relative to the snapshot instead of being coalesced away with an add that was never sent. Replay runs as the reconnecting stream's replay callback, under sendMu before the stream is published, so no chunk is sent in between. A chunk that waits longer than the ack timeout hangs the stream up so the reconnect replays it. Flush now waits for the acks, so a nil Flush means the engine committed the set. Because Flush is called before Actions by the serverless link, the session runs its receive and heartbeat loops from connect until Close; the receive loop demultiplexes acks and actions, buffering actions in an inbox that the consumer Actions attaches drains. Starting the loops and attaching the consumer add to the WaitGroup under the mutex Close takes, which removes the Actions/Close race on the WaitGroup; Actions after Close is refused. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK
| Commit: | 417f8c3 | |
|---|---|---|
| Author: | Alexander Belanger | |
| Committer: | Alexander Belanger | |
fix(grpcoperator): acknowledge action deltas and retain them until committed The Listen protocol had no delta acknowledgement. A delta the stream had accepted but the engine had not committed (stream ended before the apply, database error, rejected chunk) was dropped by the client as soon as the send returned, and a resumed worker never learned about it: the client only replayed its desired set for a fresh worker, and repeating the same AddActions was a no-op against the desired set. Listen now returns a stream of OperatorListenResponse, a oneof of the dispatcher's AssignedAction and an OperatorActionsAck. Each delta carries a positive, strictly increasing sequence; the engine applies deltas in stream order and acks a sequence once the change is committed, and an ack for N confirms every lower sequence. A delta that cannot be applied, or whose ack cannot be written, ends the stream instead of leaving it unconfirmed. The dispatcher session wraps actions into the envelope and writes acks through the same per-stream send serialisation as the fan-out. The client keeps every sent chunk until its ack arrives and replays the unacked chunks in order on every reconnect. When the reconnect did not resume the worker it replaces the queue with a snapshot of the desired set, taken under the same lock AddActions and RemoveActions use, so a removal that races the replay is queued relative to the snapshot instead of being coalesced away with an add that was never sent. Replay runs as the reconnecting stream's replay callback, under sendMu before the stream is published, so no chunk is sent in between. A chunk that waits longer than the ack timeout hangs the stream up so the reconnect replays it. Flush now waits for the acks, so a nil Flush means the engine committed the set. Because Flush is called before Actions by the serverless link, the session runs its receive and heartbeat loops from connect until Close; the receive loop demultiplexes acks and actions, buffering actions in an inbox that the consumer Actions attaches drains. Starting the loops and attaching the consumer add to the WaitGroup under the mutex Close takes, which removes the Actions/Close race on the WaitGroup; Actions after Close is refused. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK
| Commit: | 301dfab | |
|---|---|---|
| Author: | Alexander Belanger | |
feat(serverless): define the endpoint wire contract in protobuf api-contracts/v1/serverless.proto is the single source for every message the serverless operator exchanges with an endpoint: the healthcheck request and response, the non-durable trigger request and its error override, and the durable websocket frame (a oneof of first, request, response, error and done). Messages travel as protojson, so the Go operator and the TypeScript package read and write the same encoding from the same definitions. The Go bindings are generated into internal/services/shared/proto/v1 with the same dual proto path as operator.proto, since both reference .AssignedAction from the package-less dispatcher.proto. The TypeScript bindings live in sdks/typescript-serverless, the future @hatchet-dev/serverless package, which for now holds only the generation script (ts-proto pinned to the SDK's version and options) and its output under src/generated/proto. The SDK's own protoc script skips serverless.proto the way it skips operator.proto. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK
| Commit: | 7cb0b07 | |
|---|---|---|
| Author: | Alexander Belanger | |
| Committer: | abelanger5 | |
feat(engine): gRPC operator service for out-of-process operators Adds v1.OperatorService so an operator can run outside the engine and talk to it over gRPC with a tenant-scoped token: - Register (unary) upserts the operator by (tenant, name, kind = GRPC) and creates the worker for the connection, or resumes a previous worker of the same operator. Registration carries no workflows and no actions. - Listen (bidi) activates the worker for the lifetime of the stream: the first message is start{worker_id}, then heartbeats and action deltas flow in and AssignedAction messages flow out straight from the dispatcher fan-out. Deltas are applied incrementally (at most 1000 ids per message) and the scheduler is notified at most once per second while they arrive. - SendStepActionEvent and DurableTask delegate to the dispatcher; the durable stream is wrapped so the first register_worker message is checked for worker ownership before the dispatcher sees it. - Every RPC after Register, including Listen, carries hatchet-operator-id metadata validated against the token's tenant through a cached lookup. Repository: CreateWorkerOpts.OperatorId flows into the existing CreateWorker query and skips WORKER / WORKER_SLOT metering, since operator workers are excluded from those limit counts. AddWorkerActions / RemoveWorkerActions apply deltas in one transaction under the worker row lock and maintain the action hash as the XOR of per-action sha256 digests, so equal sets hash equal in any order and a delta costs O(chunk). The proto keeps importing the package-less dispatcher.proto for the shared action types (no SDK regeneration). SERVER_GRPC_OPERATORS_ENABLED (default false) gates service registration. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK
| Commit: | 44864ac | |
|---|---|---|
| Author: | Alexander Belanger | |
| Committer: | Alexander Belanger | |
feat(engine): gRPC operator service for out-of-process operators Adds v1.OperatorService so an operator can run outside the engine and talk to it over gRPC with a tenant-scoped token: - Register (unary) upserts the operator by (tenant, name, kind = GRPC) and creates the worker for the connection, or resumes a previous worker of the same operator. Registration carries no workflows and no actions. - Listen (bidi) activates the worker for the lifetime of the stream: the first message is start{worker_id}, then heartbeats and action deltas flow in and AssignedAction messages flow out straight from the dispatcher fan-out. Deltas are applied incrementally (at most 1000 ids per message) and the scheduler is notified at most once per second while they arrive. - SendStepActionEvent and DurableTask delegate to the dispatcher; the durable stream is wrapped so the first register_worker message is checked for worker ownership before the dispatcher sees it. - Every RPC after Register, including Listen, carries hatchet-operator-id metadata validated against the token's tenant through a cached lookup. Repository: CreateWorkerOpts.OperatorId flows into the existing CreateWorker query and skips WORKER / WORKER_SLOT metering, since operator workers are excluded from those limit counts. AddWorkerActions / RemoveWorkerActions apply deltas in one transaction under the worker row lock and maintain the action hash as the XOR of per-action sha256 digests, so equal sets hash equal in any order and a delta costs O(chunk). The proto keeps importing the package-less dispatcher.proto for the shared action types (no SDK regeneration). SERVER_GRPC_OPERATORS_ENABLED (default false) gates service registration. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_014gXd9VgU8FwoN3DijfprEK
| Commit: | 5ba3d7a | |
|---|---|---|
| Author: | mrkaye97 | |
feat: add cancellation reason
| Commit: | f188d07 | |
|---|---|---|
| Author: | Julius Park | |
initial commit
| Commit: | 67b18fc | |
|---|---|---|
| Author: | abelanger5 | |
| Committer: | GitHub | |
feat(engine): dynamic per-group max runs via CEL expression (#4863)
| Commit: | d6a6d5f | |
|---|---|---|
| Author: | Alexander Belanger | |
| Committer: | Alexander Belanger | |
feat: dynamic per-group max runs via CEL expression Adds Concurrency.max_runs_expression: a CEL expression over task input that computes the max runs for that task's concurrency group, so different groups (e.g. pricing tiers) get different limits from one strategy. Supported only by the in-memory concurrency index; dynamic strategies always take that path. The expression evaluates at task-insert time next to the key expression (the scheduler never sees task input) and the resulting value rides on the concurrency slot through the insert triggers, WAL payloads, and hydration. A group's effective limit is the value evaluated for its most recently created task: the in-memory subQueue applies INSERT observations guarded by task_inserted_at, so replays and retry re-inserts of older tasks cannot regress a newer task's limit. The subQueue's begin/rollback scope snapshots the limit so a failed flush restores it with the membership. Evaluation failures (non-integer or non-positive results) fail the task like key-expression failures. Registration statically rejects expressions that can never return an integer and workflow-level entries on the old parent/child DAG path. Tenant-scoped strategies carry the expression through the definition-sync trigger, checksum stripping, and strategyDiffers-driven manager rebuilds. All schema changes live in their own migration; SDK surfaces ship on a separate branch (this change carries only the regenerated contract bindings). Covered by unit tests on the subQueue guard and decide behavior under raised/lowered limits, plus integration tests through the real triggers and outbox, including the replay-cannot-regress path and chained-slot hand-off. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_0174j12vY8XiKEkkoZqNnC89
| Commit: | 67461cf | |
|---|---|---|
| Author: | abelanger5 | |
| Committer: | GitHub | |
feat(engine): shared concurrency strategies (#4845)
| Commit: | 02796b0 | |
|---|---|---|
| Author: | Alexander Belanger | |
| Committer: | Alexander Belanger | |
feat: dynamic per-group max runs via CEL expression Adds Concurrency.max_runs_expression: a CEL expression over task input that computes the max runs for that task's concurrency group, so different groups (e.g. pricing tiers) get different limits from one strategy. Supported only by the in-memory concurrency index; dynamic strategies always take that path. The expression evaluates at task-insert time next to the key expression (the scheduler never sees task input) and the resulting value rides on the concurrency slot through the insert triggers, WAL payloads, and hydration. A group's effective limit is the value evaluated for its most recently created task: the in-memory subQueue applies INSERT observations guarded by task_inserted_at, so replays and retry re-inserts of older tasks cannot regress a newer task's limit. The subQueue's begin/rollback scope snapshots the limit so a failed flush restores it with the membership. Evaluation failures (non-integer or non-positive results) fail the task like key-expression failures. Registration statically rejects expressions that can never return an integer and workflow-level entries on the old parent/child DAG path. Tenant-scoped strategies carry the expression through the definition-sync trigger, checksum stripping, and strategyDiffers-driven manager rebuilds. All schema changes live in their own migration; SDK surfaces ship on a separate branch (this change carries only the regenerated contract bindings). Covered by unit tests on the subQueue guard and decide behavior under raised/lowered limits, plus integration tests through the real triggers and outbox, including the replay-cannot-regress path and chained-slot hand-off. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_0174j12vY8XiKEkkoZqNnC89
| Commit: | 7133e73 | |
|---|---|---|
| Author: | Alexander Belanger | |
refactor: address review feedback on the backend surface - rename the proto discriminator to is_tenant_scoped and regenerate bindings - rename the unique constraint to the _uq convention - upsert a workflow's tenant strategy definitions in one statement instead of a per-definition loop (row locks are taken in name order inside the statement) - ORDER BY in the tenant stale sweep for deterministic lock acquisition - clarify the 25-hour staleness window and the non-CONCURRENTLY index build - drop a stray blank line in the proto Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_0174j12vY8XiKEkkoZqNnC89
| Commit: | fd9af62 | |
|---|---|---|
| Author: | Alexander Belanger | |
| Committer: | Alexander Belanger | |
feat: dynamic per-group max runs via CEL expression Adds Concurrency.max_runs_expression: a CEL expression over task input that computes the max runs for that task's concurrency group, so different groups (e.g. pricing tiers) get different limits from one strategy. Supported only by the in-memory concurrency index; dynamic strategies always take that path. The expression evaluates at task-insert time next to the key expression (the scheduler never sees task input) and the resulting value rides on the concurrency slot through the insert triggers, WAL payloads, and hydration. A group's effective limit is the value evaluated for its most recently created task: the in-memory subQueue applies INSERT observations guarded by task_inserted_at, so replays and retry re-inserts of older tasks cannot regress a newer task's limit. The subQueue's begin/rollback scope snapshots the limit so a failed flush restores it with the membership. Evaluation failures (non-integer or non-positive results) fail the task like key-expression failures. Registration statically rejects expressions that can never return an integer and workflow-level entries on the old parent/child DAG path. Tenant-scoped strategies carry the expression through the definition-sync trigger, checksum stripping, and strategyDiffers-driven manager rebuilds. Covered by unit tests on the subQueue guard and decide behavior under raised/lowered limits, integration tests through the real triggers and outbox (including the replay-cannot-regress path), and Go e2e tests against a live engine. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_0174j12vY8XiKEkkoZqNnC89
| Commit: | ccf6d8a | |
|---|---|---|
| Author: | Alexander Belanger | |
feat: explicit tenant_scoped flag and ordering-conflict rejection Concurrency entries gain a tenant_scoped boolean as the explicit discriminator: name alone no longer implies tenant scope (leaving names usable for workflow-scoped entries later), and expression is required on every entry. Reference-only entries are gone; every task declaring a tenant-scoped strategy carries its full definition, upserted in place. Registration now rejects chains that order the same tenant-scoped strategies inconsistently, since tasks hold earlier slots in their chain while queued on later ones and conflicting orders can deadlock. The per-step ordered chains already exist implicitly as the referencing rows' creation order, so the check aggregates them per latest workflow version (excluding the workflow being re-registered, so reorders stay possible), combines them with the incoming chains, and runs cycle detection over the pairwise precedence graph, which also catches cycles spread across more than two chains. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_0174j12vY8XiKEkkoZqNnC89
| Commit: | c6dd784 | |
|---|---|---|
| Author: | Alexander Belanger | |
refactor: tenant-scoped concurrency as ordered Concurrency entries Replace the dedicated shared-concurrency proto surface with an optional name on the Concurrency message: entries in CreateTaskOpts.concurrency and CreateWorkflowVersionRequest.concurrency_arr are processed in array order, so a tenant-scoped entry may come before or after a workflow-scoped one. A named entry with an expression defines (or updates in place) the tenant strategy as part of registration; a name-only entry references one defined elsewhere. The shared_concurrency field, shared_concurrency_defs, and the SharedConcurrencyDef message are removed. Ordering: strategy rows are created in array order, and the per-step chain follows creation order (ascending row id); rows referencing a tenant strategy resolve to the tenant strategy id only after sorting, so the tenant strategy's own id never perturbs the declared order. Workflow-level tenant-scoped entries are supported on the new DAGs-as-durable-tasks path by attaching them to the DAG orchestrator task's concurrency (workflow-scoped entries keep the parent fan-out, which always gates first); single-task workflows inherit them via the existing merge. On the old parent/child path they are rejected loudly. The workflow version checksum hashes tenant-scoped entries by name and position only, so re-registering an unchanged workflow with a changed definition reuses the version while still updating the strategy in place (definitions are upserted before the checksum early-return). SDKs: types.Concurrency (Go) and the Concurrency options (TS) gain an optional name; Python's SharedConcurrency now maps onto the Concurrency proto. All three SDKs pass a single ordered concurrency list through. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_0174j12vY8XiKEkkoZqNnC89
| Commit: | be0aede | |
|---|---|---|
| Author: | Alexander Belanger | |
| Committer: | Alexander Belanger | |
feat: tenant-scoped shared concurrency strategies Allow v1 concurrency strategies that are not workflow-scoped: a shared strategy is defined per tenant (unique name) in the new v1_tenant_concurrency table and referenced by tasks across different workflows, so they consume the same concurrency limit. Only the in-memory concurrency index supports these strategies. Schema: - new v1_tenant_concurrency table drawing ids from v1_step_concurrency's identity sequence, so strategy ids stay unique across both tables (slots, outbox topics, leases, and advisory locks key on the bare id); the borrow is documented at both definition sites since Postgres records no dependency for it - v1_step_concurrency keeps its shape and gains a nullable tenant_strategy_id: a row referencing a tenant strategy carries a point-in-time copy of the definition, but reads always resolve through the referenced row - the slot-insert reactivation trigger gains a tenant-strategy branch Engine: - strategy definitions travel as SharedConcurrencyDef entries on CreateWorkflowVersionRequest and are upserted inside the workflow-put transaction before steps are created; there is no standalone registration RPC - definitions are excluded from the workflow version checksum: they are tenant-level state, not workflow shape, so changing one updates the strategy in place without minting new workflow versions (the names steps reference remain part of the hash) - referencing rows never get a ConcurrencyManager: the lease listing returns the tenant strategy instead (zero-uuid workflow columns mark tenant scope in the shared descriptor type) - tenant strategies always use the in-memory index, skip the workflow-version active check, retire via a tenant stale sweep, and their managers are rebuilt when a definition changes, so re-registering with new settings propagates within one lease-poll cycle - per-step strategy lookup resolves the tenant strategy's id and definition, so task slots carry the shared id - fix a latent scheduler panic: RunConcurrencyStrategy returned a nil result for strategies with no runnable branch (e.g. NONE), which notifyAfterConcurrency dereferenced SDKs (Go, Python, TypeScript): - shared strategies are declared in task-level concurrency options; definitions collected from tasks ride in the workflow put request, and name-only entries reference strategies registered elsewhere - e2e tests in all three SDKs cover cross-workflow serialization, both directions of the mixed inline+shared case, and an overlap control Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_0174j12vY8XiKEkkoZqNnC89
| Commit: | 388746c | |
|---|---|---|
| Author: | Alexander Belanger | |
| Committer: | Alexander Belanger | |
feat: tenant-scoped shared concurrency strategies Allow v1 concurrency strategies that are not workflow-scoped: a shared strategy is registered per tenant (unique name) in the new v1_tenant_concurrency table and referenced by tasks across different workflows, so they consume the same concurrency limit. Only the in-memory concurrency index supports these strategies. Schema: - new v1_tenant_concurrency table drawing ids from v1_step_concurrency's identity sequence, so strategy ids stay unique across both tables (slots, outbox topics, leases, and advisory locks key on the bare id) - v1_step_concurrency keeps its shape and gains a nullable tenant_strategy_id: a row referencing a tenant strategy carries a point-in-time copy of the definition, but reads always resolve through the referenced row - the slot-insert reactivation trigger gains a tenant-strategy branch Engine: - referencing rows never get a ConcurrencyManager: the lease listing returns the tenant strategy instead (zero-uuid workflow columns mark tenant scope in the shared descriptor type) - tenant strategies always use the in-memory index, skip the workflow-version active check, retire via a tenant stale sweep, and their managers are rebuilt when a definition changes, so re-registering with new settings propagates within one lease-poll cycle - per-step strategy lookup resolves the tenant strategy's id and definition, so task slots carry the shared id - strategy definitions travel on CreateWorkflowVersionRequest and are upserted inside the workflow-put transaction before steps are created; the standalone PutSharedConcurrencyStrategy RPC remains for tuning limits without a deploy - fix a latent scheduler panic: RunConcurrencyStrategy returned a nil result for strategies with no runnable branch (e.g. NONE), which notifyAfterConcurrency dereferenced SDKs (Go, Python, TypeScript): - shared strategies are declared in task-level concurrency options; definitions collected from tasks ride in the workflow put request, and name-only entries reference strategies registered elsewhere - e2e tests in all three SDKs cover cross-workflow serialization, both directions of the mixed inline+shared case, and an overlap control Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_0174j12vY8XiKEkkoZqNnC89
| Commit: | e32f645 | |
|---|---|---|
| Author: | Julius Park | |
| Committer: | GitHub | |
feat: `CANCEL_QUEUED_EXCEPT_NEWEST`/`CANCEL_QUEUED_EXCEPT_OLDEST` concurrency strategies (#4793)
| Commit: | 7bccc19 | |
|---|---|---|
| Author: | Julius Park | |
initial commit
| Commit: | 208bbd1 | |
|---|---|---|
| Author: | mrkaye97 | |
| Committer: | mrkaye97 | |
fix: claude attempt at a fix
| Commit: | 914d9e4 | |
|---|---|---|
| Author: | mrkaye97 | |
chore: clean up sdk changes
| Commit: | 8b34cfb | |
|---|---|---|
| Author: | mrkaye97 | |
Merge branch 'main' into fix--durable-satisfied-callback-order-gurantees
| Commit: | e667a4e | |
|---|---|---|
| Author: | Julius Park | |
| Committer: | GitHub | |
feat: Batch flush (#3975)
| Commit: | 3c53ad3 | |
|---|---|---|
| Author: | Julius Park | |
| Committer: | Julius Park | |
Merge branch 'main' into feat--batch-flush-redux
| Commit: | 939ae47 | |
|---|---|---|
| Author: | matt | |
| Committer: | GitHub | |
feat(engine): terminal status-based idempotency key release (#4460) * feat: queue idempotency release * feat: col in db * feat: proto * feat: wire it all up * fix: query * feat: python impl * feat: e2e test * fix: improve query * fix: cron input test * fix: small sleep for flakiness * fix: test * test: extend idempotency tests * feat: other sdks * docs: add notes on idempotency strategies * docs: more status details * fix: move release later * fix: clean up query * test: python tests * fix: don't need to join to Step anymore * chore: gen * chore: gen / i hate this * chore: gen
| Commit: | 4c9f96e | |
|---|---|---|
| Author: | mrkaye97 | |
feat: python impl
| Commit: | 63496ef | |
|---|---|---|
| Author: | mrkaye97 | |
feat: wire it all up
| Commit: | 3ac5746 | |
|---|---|---|
| Author: | mrkaye97 | |
feat: proto
| Commit: | ba9baec | |
|---|---|---|
| Author: | Julius Park | |
merge
| Commit: | 34c8b92 | |
|---|---|---|
| Author: | matt | |
| Committer: | GitHub | |
feat(engine): idempotency (#4045) * fix: small optimization, no need to query if there are no claims * feat: add idempotency key expr to proto, workflow version * feat: wiring * fix: wiring to db * feat: add ttl col * chore: gen a whole bunch of python * feat: wire up idempotency config on the sdk * fix: engine wiring * fix: py type * chore: revert worker changes * feat: new example * feat: fix migration, add cols * feat: wiring * feat: wire up idempotency to the triggerTuple * feat: initial engine wiring * fix: len * fix: types * feat: add test * chore: rm unused deps * fix: use the v1 trigger endpoint on python * fix: make the python test do a real thing * feat: first pass bypassing mq fallback for idempotency * fix: remove todo, return the same id * feat: wire up existing run ids * fix: use test run id to dedupe across test runs * fix: improve test a bit * chore: lock * fix: add grpcio types * fix: drive by fix for a flaky test caused by event condition race * fix: not found behavior on the get_details method * fix: revert python changes * fix: handle collisions in the v1 ingest endpoint :facepalm: * chore: gen * chore: add back some comments to shrink diff * chore: python version, changelog * fix: use the correct tx * fix: couple query bugs * fix: couple more event changes * feat: idempotency test with events * fix: type issue * fix: tests, lint * fix: test cleanup * fix: handle deduplication in the query * fix: test, gen * fix: conflict * fix: add prefix to idempotency key claim, clean up a tiny bit * fix: migration name, add idempotency key to cel eval failure source enum * fix: pr comments * chore: gen * feat: snippets, initial doc work * chore: vibe code other sdks * chore: vendor the status proto for js * fix: go docs * fix: ruby simplification * chore: changelogs * fix: ruby ci * fix: more test fixes * chore: gen * chore: gen again * fix: tests * chore: gen * chore: changelogs * chore: lockfile * chore: lint * chore: appease the cop * chore: fix copilot comments * chore: gen * chore: migration version * chore: migration version * fix: ordering to prevent deadlocks * fix: add validation, add skeleton for handling collisions on bulk trigger * feat: add key to v1_task * feat: wiring * feat: add idempotency keys to runs, dags, etc. tables * feat: wire up sending idempotency back over the api to the dash * feat: wire up idempotency filter and column on the fe * feat: return partial success on bulk trigger via an error * feat: wrap bulk triggers * feat: add test * feat: go migration for indexes * fix: order * fix: down * fix: weird diff * fix: always create trigger writer * chore: gen other sdks * fix: queries * fix: param ordering * chore: gen * chore: prettier * fix: merge * fix: migration * chore: lint * chore: bundle * chore: gen changelogs * feat: add TTL-based idempotency config + ABC to allow for better extensibility for other strategies * chore: other sdks, gen * chore: lint * fix: rubocop * fix: rm generic * chore: schema
| Commit: | 08ffdf6 | |
|---|---|---|
| Author: | mrkaye97 | |
feat: return partial success on bulk trigger via an error
| Commit: | fed9fa4 | |
|---|---|---|
| Author: | Julius Park | |
refactor so that messages are sent in one shot
| Commit: | 94f1c92 | |
|---|---|---|
| Author: | mrkaye97 | |
Merge branch 'main' into mk/feat-idempotency
| Commit: | a441554 | |
|---|---|---|
| Author: | mrkaye97 | |
chore: rm unused proto change
The documentation is generated from this commit.
| Commit: | 02f6fd0 | |
|---|---|---|
| Author: | Gabe Ruttner | |
chore order of things
| Commit: | 115306a | |
|---|---|---|
| Author: | Gabe Ruttner | |
comments
| Commit: | e520b06 | |
|---|---|---|
| Author: | Gabe Ruttner | |
refactor: more feedback and clarified naming
| Commit: | e8d2ba7 | |
|---|---|---|
| Author: | Gabe Ruttner | |
ack
| Commit: | a1a84b4 | |
|---|---|---|
| Author: | Gabe Ruttner | |
chore: refactor rate limit stuff
| Commit: | 10301e1 | |
|---|---|---|
| Author: | Gabe Ruttner | |
feat: separate stream proto
| Commit: | 91d1343 | |
|---|---|---|
| Author: | Gabe Ruttner | |
chore: easy feedback
| Commit: | 151fd44 | |
|---|---|---|
| Author: | Gabe Ruttner | |
Merge branch 'main' into feat--hatchet-v2-grpc-api
| Commit: | ac147cb | |
|---|---|---|
| Author: | mrkaye97 | |
feat: sql, move some logic over from my old branch
| Commit: | 91549f1 | |
|---|---|---|
| Author: | Gabe Ruttner | |
fix change og
| Commit: | bd482af | |
|---|---|---|
| Author: | Gabe Ruttner | |
Merge branch 'main' into fix--durable-satisfied-callback-order-gurantees
| Commit: | 1963828 | |
|---|---|---|
| Author: | Julius Park | |
Merge branch 'main' into feat--batch-flush-redux
| Commit: | 15addec | |
|---|---|---|
| Author: | matt | |
| Committer: | GitHub | |
fix(engine): correctly raise error in durable parent when child fails (#4154) * fix: propagate error back to the sdk * fix: python sdk * feat: add is_failure and error_message to the proto * feat: engine changes for wiring failure flag + error msg through * chore: other sdks * chore: cop * chore: nullif on write * fix: feedback * feat: e2e test * feat: tests * fix: tests * chore: lint * fix: test * chore: gen * chore: changelog sync * chore: lint ugh I hate our ci * fix: add tasks to worker * chore: gen * chore: versions, changelogs * chore: bundle * fix: rework test to fail on cloud * fix: autocomplete :facepalm: * chore: black * chore: gen
The documentation is generated from this commit.
| Commit: | 2e0f991 | |
|---|---|---|
| Author: | Gabe Ruttner | |
fix: track satisfied order
| Commit: | 8d37fa8 | |
|---|---|---|
| Author: | mrkaye97 | |
fix: move graph creation to the engine
| Commit: | f29a397 | |
|---|---|---|
| Author: | mrkaye97 | |
feat: almost working e2e PoC
| Commit: | 552a47e | |
|---|---|---|
| Author: | Gabe Ruttner | |
separate register and stream
| Commit: | ab7d578 | |
|---|---|---|
| Author: | Gabe Ruttner | |
feat: proposed protos
| Commit: | 0fb5bf6 | |
|---|---|---|
| Author: | Julius Park | |
add result broadcasting
| Commit: | 0bdfc1c | |
|---|---|---|
| Author: | mrkaye97 | |
chore: vendor the status proto for js
| Commit: | 17e5cbc | |
|---|---|---|
| Author: | Julius Park | |
rename proto fields
| Commit: | 4942ad4 | |
|---|---|---|
| Author: | mrkaye97 | |
feat: wire up existing run ids
| Commit: | 564d867 | |
|---|---|---|
| Author: | Julius Park | |
| Committer: | Julius Park | |
wip merge
| Commit: | 247dd6c | |
|---|---|---|
| Author: | mrkaye97 | |
feat: add ttl col
| Commit: | 09b749e | |
|---|---|---|
| Author: | mrkaye97 | |
feat: add idempotency key expr to proto, workflow version
| Commit: | 5e1fabb | |
|---|---|---|
| Author: | matt | |
| Committer: | GitHub | |
Fix: Webhook responses, event info on context, internal fix for labels matches (#3625) * feat: add boolean flag for whether or not to return the event as the response payload * feat: add field to indicate whether or not to return the event * fix: wire up receive api * feat: more columns for event trigger idea * feat: start wiring up queries * feat: internal wiring * fix: matches, more wiring * fix: simplify query * fix: more wiring * fix: clean up dag impl for worker labels too * fix: wiring, unwind dag changes * fix: if * feat: add webhook response switch on fe * fix: casing * fix: populator * feat: add fields to assigned action for trigger * feat: send trigger data back over the api to the context * chore: gen py * feat: wire up context in py * feat: update webhooks feature client * feat: ts, go * chore: lint * fix: more wiring for replay * fix: nil handling * chore: gen * chore: attempt to fix flaky test * feat: wire up olap writes * fix: rm changes on the olap side for now * fix: improve the `was_triggered_by_event` prop * chore: gen, fix migration version * chore: versions, changelogs
| Commit: | b39548c | |
|---|---|---|
| Author: | matt | |
| Committer: | GitHub | |
Feat: Durable execution frontend work and API improvements (#3639) * feat: initial work on durable event log ui * feat: use the logging component to show the event log * remove evicted crescent since the hover state was broken and it looked awkward * fix: filter out memo events, change copy * fix: durable event log height * feat: add helper hints for run triggers too * feat: rework wait data * feat: bulk run grouping * fix: or groups, more grouping on fe * feat: simplify api more * fix: get rid of cursor-impl mini map click behavior * feat: start consolidating durable and non-durable event logs * fix: make log level an enum for log component * feat: improve color of different logs, etc * fix: show task display names for dags * refactor: move durable event list to durable events repo, update api to return external ids and display names consistently * fix: swap order of replay / cancel buttons * feat: tabs cleanup for side panel * fix: input, output, additional metadata tab sizes in the side panel * chore: lint * feat: add link to log line * fix: improve copy * chore: lint and gen * feat: label in ts * chore: gen protos * fix: default behavior * chore: versions, changelogs * chore: rm some comments * chore: docs * fix: tests * fix: one more satisfied-now fix * fix: another copilot suggestion * fix: more copilot * fix: pagination * fix: start fixing link clicks * fix: log click overflow logic * Fix: Misc frontend issues (#3633) * fix: event column truncation * feat: add popover for worker labels * fix: memoize table cols * fix: simple table scrolling and padding * fix: lint * fix: dag view * fix: populator type assert * fix: remove duration string helper * fix: durable task ids
| Commit: | d69ef31 | |
|---|---|---|
| Author: | matt | |
| Committer: | GitHub | |
Feat: Wait for event with lookback window (#3442) * feat: add scope to proto * chore: gen * feat: pass scope through to wait from sdk * feat: lookback window * feat: wire up opts * feat: add lookup query * feat: look up historical matches * fix: enable durable event log for everyone * fix: array len bug * fix: rm deduplication * fix: throw if either scope or lookback window is provided but not both * fix: simplify * fix: lint docs * fix: lint * refactor: pull satisfying matches into a separate tx * fix: forgot to commit :facepalm: * refactor: factor out logic into helper * feat: add some tests for lookbacks * fix: test cleanup, lint * fix: modify index, add scope * Revert "fix: modify index, add scope" This reverts commit b0ceb91cf36ed7dcb76ad47f81d0360e76a1b38b. * chore: gen ts * fix: start passing scope and lookback window through * fix: finish ts wiring * chore: replicate tests to ts * chore: gen everything * chore: gen again, it's neverending * fix: decrease sleep time * chore: lock * fix: commit if we return early * chore: comment * chore: try pinner redocly * fix: try overriding ajv * fix: impl consider events since in js correctly * chore: changelog * fix: max ordering, last event wins * chore: lint
| Commit: | a899c32 | |
|---|---|---|
| Author: | matt | |
| Committer: | GitHub | |
Feat: Add `workflow_run_external_id` to trigger run ack proto (#3299) * feat: add external id to runs entry * chore: gen * fix: python wiring * chore: gen ts * chore: version * fix: wiring * chore: changelog, version * chore: gen
| Commit: | 43e7466 | |
|---|---|---|
| Author: | matt | |
| Committer: | GitHub | |
Feat: Durable Execution Revamp (#2954) * Feat: Durable Execution Revamp --------- Co-authored-by: Gabe Ruttner <gabriel.ruttner@gmail.com>
| Commit: | 583bb6f | |
|---|---|---|
| Author: | Gabe Ruttner | |
| Committer: | GitHub | |
Feat--unify-running-and-evicted (#3248) * feat: unify running and evicted counts with filter * chore: gen * chore: gen lint * fix: opsies... * chore: gen * chore: feedback
| Commit: | 468de0c | |
|---|---|---|
| Author: | Gabe Ruttner | |
| Committer: | GitHub | |
feat: handle server initiated eviction on workers (#3200) * refactor: ts in sync with python * fix: waits respect abort signal * chore: gen, lint * fix: update key extraction logic in durable workflow example * fix: flake * fix: racey test * chore: gen * feat: send eviction to engine * feat: handle eviction on workers * fix: tests and invocation scoped eviction * fix: cleanup scope * refactor: cleanup how we do action keys in ts * chore: gen, lint * fix: test * chore: lint * fix: mock * fix: cleanup * chore: lint * chore: copilot feedback * fix: reconnect durable listener * chore: feedback * fix: bump worker status on wait and reconnect * refactor: rename eviction notice * fix: return error on query failure * chore: lint * chore: gen * refactor: remove early completion * Revert "refactor: remove early completion" This reverts commit f756ce37b04f4eaf02d8e74615402b773e05049e. * chore: remove unused param
| Commit: | 1cd9660 | |
|---|---|---|
| Author: | Gabe Ruttner | |
| Committer: | GitHub | |
rip: remove unneeded durable event log update (#3186) * rip: update * refactor: start cleaning up proto defs * refactor: finish cleaning up proto definitions * fix: rm the kind * refactor: rewire the server to have different methods for the different paths * refactor: more intermediate work * fix: variables * fix: wire the kind through * refactor: get it to compile * chore: start fixing python * fix: first pass at fixing python * chore: rm unused sql * fix: rm invocation count, rework some logic * fix: alias (why does sqlc not catch this) * fix: panics * fix: add faster timeout to durable spawn test * fix: task id bug * refactor: more cleanup of types * refactor: rm stale entries logic * fix: rework getOrCreate logic * fix: clean up a bunch more unneeded stuff * fix: bug * fix: python code * fix: dag matches bug * fix: add parent to dag to make it more broken, add timeout * fix: more involved tests * fix: dag waits * fix: tests --------- Co-authored-by: mrkaye97 <mrkaye97@gmail.com>
| Commit: | 12bf751 | |
|---|---|---|
| Author: | Gabe Ruttner | |
| Committer: | GitHub | |
feat: durable-bulk-spawn (#3173) * feat: durable bulk spawn * feat: bulk db operations * chore: lint * chore: generate * fix: sqlc hack * cleanup * chore: lint * refactor: remove duplicated code * refactor: simplify input * ops: rules * fix: add missing workflows to worker * chore: add todo * refactor: single execution path * comment * refactor: remove extra class * chore: lint * chore: lint * refactor: batch ack * fix: handle empty refs * chore: feedback * fix: durable output * chore: lint * chore: lint * chore: test fix and gen * revert: actually use bulk * refactor: simplify path * chore: empty commit Made-with: Cursor
| Commit: | 1bffb66 | |
|---|---|---|
| Author: | matt | |
| Committer: | GitHub | |
Feat: Branching off branches (#3150) * feat: add branch point table * chore: gen * feat: id for ordering * feat: check for `isBranchPoint` and handle branching * feat: wire up branching over the api * fix: api, gen * fix: gen * fix: branch * feat: remove duped parent node and branch ids * fix: rm branch count, dupe of latest id * refactor: resolve naming * fix: base case * fix: test, + it's literally always a caching issue omg * fix: docs * chore: lint * refactor: make branch resolution more efficient * feat: stable sort, add a bunch of tests * fix: confusing naming * fix: naming * chore: gen * fix: update * fix: failing test --------- Co-authored-by: Gabe Ruttner <gabriel.ruttner@gmail.com>
| Commit: | a6f5fb6 | |
|---|---|---|
| Author: | Gabe Ruttner | |
| Committer: | GitHub | |
chore: rename durable fork to branch (#3174)
| Commit: | 5b58eed | |
|---|---|---|
| Author: | Gabe Ruttner | |
| Committer: | GitHub | |
tests: compat testing (#3144) * tests: compat testing * fix: new engine old sdk * feat: ruby ts * fix: namespaced tests, conditions * chore: lint ts * fix: python compat * chore: lint * fix: warn on version mismatch * chore: lint * fix: child spawn * refactor: address review * chore: lint * chore: lint * chore: lint * refactor: address reviews * chore: lint * fix: typeguard * tests: docker-compose * fix: docker compose * fix: labels? * chore: skip new tests * fix: backwards compat * fix: grpc proxy
| Commit: | e29459f | |
|---|---|---|
| Author: | mrkaye97 | |
chore: merge main
| Commit: | 6c29e48 | |
|---|---|---|
| Author: | matt | |
| Committer: | GitHub | |
Feat: Dynamic worker label assign (#3137) * feat: initial wiring work on desired labels * feat: initial wiring * chore: gen python * fix: use the whole desired label thing instead * fix: more wiring, improve types * fix: sql type * fix: len check * chore: gen python * fix: initial plural label work * fix: store the labels properly on the task * fix: skip cache on override * fix: bug * fix: scoping bug whoops * chore: lint * fix: send labels back over the api correctly * feat: python test * fix: lint * fix: comment * fix: override * fix: namespaces, ugh * fix: no need for error here * chore: version * feat: ruby, go, ts * feat: versions * fix: appease the rubocop * chore: lint * chore: bundle install * fix: tests * chore: lint * chore: lint more * fix: ts test * fix: rb * chore: gen * chore: reset gemfile * chore: reset changelog * fix: pgroup * fix: tests, part i * Revert "chore: reset changelog" This reverts commit b63bf7d3e56d28569dd0344059f5e89c20aa8cee. * Revert "chore: reset gemfile" This reverts commit bb848bb6f0048587467141fb3fc24da0ed81a316. * fix: go -> golang mapping hack * fix: go enums * fix: appease the cop * fix: namespace * chore: gen
| Commit: | 77f769b | |
|---|---|---|
| Author: | matt | |
| Committer: | GitHub | |
Feat: Durable Memoization (#3112) * feat: initial engine / db work * chore: gen python * feat: initial python work * feat: wiring * feat: initial test * fix: scope memo key to the run id (we might not even need this) * chore: gen, docstring * fix: mandatory result type * fix: docs * fix: log a warning and don't cache * fix: add test for non-unique keys, fix some bugs * chore: lint * chore: naming * fix: more naming * fix: comment * chore: gen * fix: docs * fix: naming, ugh * chore: gen * chore: gen * fix: always send event log entry and get ack * fix: initial rework of memo to only use stream * fix: start reworking python side * fix: finish wiring up memo put * chore: fmt * chore: comment for monday * fix: start reworking signature * fix: docstring * fix: union type for send event to improve typing a bunch * fix: memo key * chore: gen * fix: rm unused query * fix: comments
| Commit: | 2d57b67 | |
|---|---|---|
| Author: | Gabe Ruttner | |
| Committer: | GitHub | |
Feat durable olap refactor (#3115) * chore: lint * feat: counting and partitioning * feat: add reason field to DurableTaskEvictInvocationRequest and update eviction handling * fix: eviction durable execution race * chore: generate * refactor: simplified migration * refactor: address review * refactor: analyze parent tables in migraiton * fix: migration * fix: remove no txn * fix: one statement * fix: we do infact need no transaction * add down/up/down to the online migration test * fix: or multiple statements * fix: two migrations... * chore: rm old migraiton * chore: generate * chore: feedback * fix: idempotent migration * refactor: update assertions in durable tests and clean up imports in cache.py * revert: migraiton * chore: wrap down
| Commit: | daff28d | |
|---|---|---|
| Author: | Gabe Ruttner | |
| Committer: | GitHub | |
Feat: durable eviction take 2 (#3075) * feat: simplified eviction feature * fix: assign new worker id * test: shorter sleep * fix: completion race on same worker * chore: address todo * chore: lint * chore: generate * fix: n+1 queries * refactor: WasEvicted bool * feat: evicted state * chore: generate * fix: map status * fix: update PendingCallback structure to include InvocationCount * revert: comment * feat: add support for EVICTED status in waterfall component and metrics display * fix: implicit eviction * chore: readable cte * refactor: queued bool * refactor: rename eviction_policy * fix: aio only * chore: example return type * fix: map * feat: eviction error cases * refactor: change external ID maps to use UUID type * chore: feedback, cleanup * tests: additional cases * chore: generate * chore: lint * chore: lint generate * chore: clean up comments to make matt happy * refactor: more feedback * chore: add TODO for worker state reconciliation and clean up comments in eviction policy * tests: fix * chore: gen * test: increase ruby timeout... * fix: invocation count * fix: test cases * fix: stale log entry * chore: lint * revert: durable tests to use time.time * chore: lint
| Commit: | 401f9dc | |
|---|---|---|
| Author: | mrkaye97 | |
feat: bulk spawn protos
| Commit: | 28c38ef | |
|---|---|---|
| Author: | Gabe Ruttner | |
feat: simplified eviction feature
| Commit: | 6f3f6e0 | |
|---|---|---|
| Author: | matt | |
| Committer: | GitHub | |
Feat: Replay as new (or from a node) (#3055) * feat: new messages for reset * chore: gen python * feat: reset scaffolding * feat: initial work * feat: initial e2e wiring of resetting from a specific node * fix: add branch to pk * fix: wire up branches * fix: add branch to awaited entry * feat: start wiring up reset api * fix: colname * fix: add branch id more places * fix: some bugs * fix: replay * fix: replay, simplify * feat: add parent branch id * fix: start reworking parent nodes and branches * fix: parent branch wiring * fix: start fixing some bugs * fix: parent branch bug * fix: advisory lock for locking the log file to prevent concurrent modification * fix: move claude.md ignore path * fix: remove eager replays of events * fix: rm cruft * fix: cleanup more params and such * fix: return type * fix: comment * fix: comments * fix: comment * chore: gen * chore: gen * fix: decrease sleep time * chore: gen again * fix: add invocation count on event log entries, make it int32, fix toInt * fix: more wiring * chore: gen, simplify * fix: lint * fix: more zero values, I hate Go * feat: add `is_durable` to v1_task * feat: initial work wiring up dispatcher to increment log entry invocation counts * feat: wire up assigned action * fix: property * fix: send is durable through to the engine * fix: more invoc count wiring * fix: node resetting * fix:revert * fix: import * chore: gen * fix: reset -> fork * fix: rm a bunch of dead code * fix: api * fix: repo method * fix: log file locking using `FOR UPDATE` + atomic compare-and-set update * fix: move to shared repo * feat: increment invocation count on the scheduler * fix: naming * fix: make test more reliable * fix: props * fix: node id reset
| Commit: | 7e3e3b8 | |
|---|---|---|
| Author: | matt | |
| Committer: | GitHub | |
Feat: Non-determinism errors (#3041) * fix: retrieve payloads in bulk * fix: hash -> idempotency key * feat: initial hashing work * feat: check idempotency key if entry exists * fix: panic * feat: initial work on custom error for non-determinism * fix: handle nondeterminism error properly * feat: add error response, pub message to task controller * chore: lint * feat: add node id field to error proto * chore: rm a bunch of unhelpful cancellation logs * fix: conflict issues * fix: rm another log * fix: send node id properly * fix: improve what we hash * fix: improve error handling * fix: python issues * fix: don't hash or group id * fix: rm print * feat: add python test * fix: add timeout * fix: improve handling of non determinism error * fix: propagate node id through * fix: types, test * fix: make serializable * fix: no need to cancel internally anymore * fix: hide another internal log * fix: add link to docs * fix: copilot * fix: use sha256 * fix: test cleanup * fix: add error type enum * fix: handle exceptions on the worker * fix: clean up a bunch of cursor imports * fix: cursor docstring formatting * fix: simplify idempotency key func * fix: add back cancellation logs * feat: tests for idempotency keys * fix: add a couple more for priority and metadata * chore: gen * fix: python reconnect * fix: noisy error * fix: improve log * fix: don't run durable listener if no durable tasks are registered * fix: non-null idempotency keys
| Commit: | eaf6bba | |
|---|---|---|
| Author: | matt | |
| Committer: | GitHub | |
Refactor: Remove separate callback table (#3045) * fix: remove callback table * fix: type * fix: type * fix: wiring everything up * fix: result payload for replays * chore: lint * feat: set fillfactor * chore: gen * fix: simplify v1 match changes * fix: simplify v1 match wiring * fix: rm print line * fix: some more wiring * fix: wiring * chore: comments * chore: gen, proto naming * chore: comments * fix: rm comment * fix: broken listener
| Commit: | 00b875e | |
|---|---|---|
| Author: | Mohammed Nafees | |
remove useless proto