Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
45 commits
Select commit Hold shift + click to select a range
5627db1
feat: WebSocket support and transport-agnostic event indexer
jleeh Apr 16, 2026
94e412e
fix(sharedState): keep-latest pub/sub delivery + ctx-deadline-aware r…
jleeh Apr 18, 2026
cddeb01
fix(upstream): plug four upstream-instance leak paths
jleeh Apr 29, 2026
bf9f50b
fix(util): pooledGzipReadCloser closes its source on Close
jleeh Apr 30, 2026
51faede
fix(networks): clamp regressions in EvmHighestBlockNumber monotonic g…
jleeh May 8, 2026
5a46e14
fix(failover): close tracker-signal gap and add per-request fallback …
May 22, 2026
6fb6b51
fix(cache): preserve "latest"/"finalized" tags for state-read methods
May 28, 2026
9372a10
fix(eth_sendRawTransaction): treat replay-attack rejections as idempo…
jleeh Jun 1, 2026
454c54d
test: adapt failover/selection tests to rebased upstream API
jleeh Jun 1, 2026
3bbc290
fix(websocket): self-heal wedged upstream WS connections and surface …
snowkide Jun 12, 2026
faf6648
test(websocket): end-to-end regression for ungraceful upstream death
snowkide Jun 12, 2026
1886390
refactor(websocket): slim to core fixes per review; fix races
snowkide Jun 12, 2026
ccfbfda
fix(failsafe): release half-open permit on ignored outcome (breaker w…
snowkide Jun 24, 2026
a88e68a
Merge pull request #3 from linkpoolio/fix/breaker-halfopen-permit-leak
snowkide Jun 30, 2026
cbbb538
fix(initializer): keep retrying recoverable tasks when a sibling task…
shpookas Jul 20, 2026
95719d0
fix(initializer): bound each auto-retry round by TaskTimeout (hung ta…
shpookas Jul 20, 2026
42f6842
fix(initializer): stop the retry loop before taking tasksMu in Stop (…
shpookas Jul 20, 2026
1be9510
fix(ws): advance network latest tip before newHeads fan-out
jleeh Jul 22, 2026
0d8003e
docs: drop consumer-specific wording from WS tip comments
jleeh Jul 22, 2026
c1c3404
fix(initializer): reap hung tasks, fix State/Wait tip races
jleeh Jul 22, 2026
defe2e4
Merge pull request #5 from linkpoolio/fix/initializer-retry-wedge
snowkide Jul 22, 2026
845ccb8
Merge pull request #10 from linkpoolio/fix/ws-tip-advances-network-la…
snowkide Jul 22, 2026
8f4a725
fix(ws): serve cached newHeads header when HTTP latest lags tip
jleeh Jul 22, 2026
eb99874
fix(ws): route tip reads to upstreams that already have the block
jleeh Jul 22, 2026
961d3f9
fix(ws): drop parallel tip-source id; use poller partition for latest
jleeh Jul 22, 2026
d142444
fix(ws): pin tip re-fetch to EvmLeaderUpstream
jleeh Jul 22, 2026
0b81d3e
fix(ws): bypass poller debounce when tip gate needs a fresh head
jleeh Jul 23, 2026
8d21290
fix(ws): refuse stale latest when tip re-fetch misses TipHW
jleeh Jul 23, 2026
0387ea2
fix(ws): sync TipHW publish and refresh on HTTP false-negative
jleeh Jul 23, 2026
28cd7f9
fix(ws): serve cached newHeads header when tip re-fetch misses TipHW
jleeh Jul 23, 2026
53456da
fix(ws): stop TipHW inflation from fallback WS; escape on emptyish miss
jleeh Jul 23, 2026
9ba72db
fix(ws): drop cached newHeads header serve for HTTP latest
jleeh Jul 23, 2026
f78796d
fix(ws): always advance TipHW for fan-out heads, keep emptyish escape
jleeh Jul 23, 2026
08dd03b
Merge pull request #11 from linkpoolio/fix/ws-tip-http-floor-from-cac…
snowkide Jul 24, 2026
4af71ff
fix(ws): pin near-tip eth_getBlockByNumber to EvmLeaderUpstream
shpookas Jul 31, 2026
80aad9e
fix(ws): slim near-tip getBlock pin to inline Forward path
shpookas Jul 31, 2026
009958d
fix(ws): restore near-tip pin helper and laggingHits test
shpookas Jul 31, 2026
9e39885
Merge pull request #14 from linkpoolio/fix/pin-near-tip-getblock-to-e…
shpookas Jul 31, 2026
e0b2c49
chore(ci): add workflow_dispatch publish to docker.io/linkpool/erpc
shpookas Aug 3, 2026
28493f4
chore(ci): publish experiment tags to ghcr.io instead of Docker Hub
shpookas Aug 4, 2026
1384e12
chore(ci): publish erpc images to ghcr.io/linkpoolio/docker-images/erpc
shpookas Aug 4, 2026
24e56ff
fix(evm): skip lagging upstreams for eth_call(latest) and mirror head…
1marcghannam Sep 15, 2026
0c96ef1
test(evm): encode Priority Pool totalQueued incident in eth_call lag …
1marcghannam Sep 15, 2026
4616381
feat(evm): make latest-state lag gate configurable via evm.maxLatestS…
1marcghannam Sep 15, 2026
5e8fc97
chore(ci): pin pnpm 9.15.9 for alpine musl GHCR builds
snowkide Sep 17, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
88 changes: 88 additions & 0 deletions .github/workflows/ghcr-publish.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,88 @@
# Build and push ghcr.io/linkpoolio/docker-images/erpc:<tag> from an arbitrary git ref.
# Uses the existing docker-images GHCR namespace (same registry path family as other LinkPool images).
# Auth: GITHUB_TOKEN with packages:write. May require package Actions access for linkpoolio/erpc
# if the package is owned/linked to docker-images.
#
# Run: Actions → "Publish to GHCR" → workflow_dispatch
# tag: ws-pin-near-tip-test9
# ref: feat/websocket-support

name: Publish to GHCR

concurrency:
group: ghcr-publish-${{ github.event.inputs.tag }}
cancel-in-progress: false

on:
workflow_dispatch:
inputs:
tag:
description: "Image tag to push (e.g. ws-pin-near-tip-test9)"
required: true
type: string
ref:
description: "Git ref to build (branch, tag, or SHA)"
required: true
default: "feat/websocket-support"
type: string
platforms:
description: "Target platforms"
required: false
default: "linux/amd64"
type: string

permissions:
contents: read
packages: write

jobs:
publish:
runs-on: ubuntu-24.04
timeout-minutes: 60
steps:
- name: Checkout
uses: actions/checkout@v4
with:
ref: ${{ inputs.ref }}
fetch-depth: 1

- name: Set build metadata
id: meta
run: |
set -euo pipefail
SHA="$(git rev-parse HEAD)"
SHORT="$(git rev-parse --short HEAD)"
IMAGE="ghcr.io/linkpoolio/docker-images/erpc"
echo "sha=${SHA}" >> "$GITHUB_OUTPUT"
echo "short_sha=${SHORT}" >> "$GITHUB_OUTPUT"
echo "image=${IMAGE}" >> "$GITHUB_OUTPUT"
echo "Building ref=${{ inputs.ref }} sha=${SHORT} → ${IMAGE}:${{ inputs.tag }}"

- name: Set up Docker Buildx
uses: docker/setup-buildx-action@v3

- name: Login to GitHub Container Registry
uses: docker/login-action@v3
with:
registry: ghcr.io
username: ${{ github.repository_owner }}
password: ${{ secrets.GITHUB_TOKEN }}

- name: Build and push
uses: docker/build-push-action@v6
with:
context: .
file: ./Dockerfile
push: true
platforms: ${{ inputs.platforms }}
tags: |
${{ steps.meta.outputs.image }}:${{ inputs.tag }}
build-args: |
VERSION=${{ inputs.tag }}
COMMIT_SHA=${{ steps.meta.outputs.short_sha }}
provenance: false
sbom: false

- name: Print digest
run: |
docker buildx imagetools inspect "${{ steps.meta.outputs.image }}:${{ inputs.tag }}"
8 changes: 4 additions & 4 deletions Dockerfile
Original file line number Diff line number Diff line change
Expand Up @@ -35,20 +35,20 @@ RUN go build -v -ldflags="$LDFLAGS" -a -installsuffix cgo -o erpc-server ./cmd/e

# Global typescript related image
FROM node:20-alpine@sha256:09e2b3d9726018aecf269bd35325f46bf75046a643a66d28360ec71132750ec8 AS ts-core
RUN npm install -g pnpm
RUN npm install -g pnpm@9.15.9

# Stage where we will install dev dependencies + compile sdk
FROM ts-core AS ts-dev
RUN mkdir -p /temp/dev/typescript
RUN npm install -g pnpm
RUN npm install -g pnpm@9.15.9

# Copy only the TypeScript package files
COPY typescript/config /temp/dev/typescript/config
COPY pnpm* /temp/dev/
COPY package.json /temp/dev/package.json

# Install everything and build
RUN --mount=type=cache,id=pnpm,target=/pnpm/store cd /temp/dev && pnpm install --frozen-lockfile
RUN --mount=type=cache,id=pnpm,target=/pnpm/store cd /temp/dev && pnpm install --frozen-lockfile --config.manage-package-manager-versions=false
RUN cd /temp/dev && pnpm build

# Stage where we will install prod dependencies only
Expand All @@ -60,7 +60,7 @@ COPY pnpm* /temp/prod/
COPY package.json /temp/prod/package.json

# Install every prod dependencies
RUN --mount=type=cache,id=pnpm,target=/pnpm/store cd /temp/prod && pnpm install --prod --frozen-lockfile
RUN --mount=type=cache,id=pnpm,target=/pnpm/store cd /temp/prod && pnpm install --prod --frozen-lockfile --config.manage-package-manager-versions=false

# Create symlink stage (for backwards compatibility with earlier image file structure)
FROM alpine:latest@sha256:25109184c71bdad752c8312a8623239686a9a2071e8825f20acb8f2198c3f659 AS symlink
Expand Down
95 changes: 91 additions & 4 deletions architecture/evm/block_ref.go
Original file line number Diff line number Diff line change
Expand Up @@ -75,10 +75,10 @@ func ExtractBlockReferenceFromRequest(ctx context.Context, r *common.NormalizedR
// In case of "*" since it means any block, we can still augment it from response ref, because during cache.Get()
// we'll be using reverse index (i.e. ignoring ref), but after reorg invalidation is added a specific block ref is useful.
//
// TODO An ideal version stores the data for all eth_getBlockByNumber(latest) and eth_getBlockByNumber(blockNumber),
// and eth_getBlockByNumber(blockHash) where blockNumber/blockHash are the actual values returned in the response.
// So that if user gets the latest block, then cache is populated for when they provide that specific block as well.
// When implementing that feature remember that CacheHash() must be calculated separately for each number/hash combo.
// For moving-tag requests ("latest", "finalized", "safe"), the
// cache layer calls ResolveCacheBlockRef instead, which resolves
// the tag to a concrete block number so each tip advance gets
// its own cache key.
blockRef = br
}
if bn > 0 {
Expand Down Expand Up @@ -114,6 +114,93 @@ func ExtractBlockReferenceFromRequest(ctx context.Context, r *common.NormalizedR
return blockRef, blockNumber, nil
}

// ResolveCacheBlockRef returns the block reference the cache layer should use
// when keying an eth_getBlockByNumber("latest") response (and other moving
// tags). Regular ExtractBlockReferenceFromRequest preserves the literal tag
// string ("latest") as blockRef so the cache hits on repeat tag queries —
// but that makes every request within the TTL window return the same pinned
// response regardless of chain progression (see the bug fixed alongside this
// helper: stale "latest" responses served from cache until TTL expiry, with
// enforceHighestBlock explicitly skipping cached responses).
//
// This helper substitutes the tag with a concrete block number so each tip
// advance is a distinct cache key: on WRITE we use the response's own block
// number (definitive answer for what the cached payload represents); on READ
// we consult the network's tip tracker (EvmHighestLatestBlockNumber, which
// aggregates max over upstream pollers and the cross-pod shared counter) to
// decide which block we'd be asking for *right now*. Within a single tip
// the key is stable and concurrent "latest" queries coalesce onto one cached
// entry; across tip advances the key changes and the next request forwards
// upstream.
//
// The function does NOT mutate the request's EvmBlockRef — the original
// "latest" tag is preserved on the request so downstream finality computation
// and other tag-aware logic keeps working.
//
// Fallback: if the tag can't be resolved to a concrete block number (no
// response, no network attached to the request, or the tracker hasn't seen
// a block yet), the original tag is returned and the cache key stays
// tag-literal — same as prior behaviour. That path should be rare in
// production since every normal HTTP request has a Network and an upstream
// response by the SET stage.
func ResolveCacheBlockRef(ctx context.Context, req *common.NormalizedRequest, resp *common.NormalizedResponse) (string, int64, error) {
blockRef, blockNumber, err := ExtractBlockReferenceFromRequest(ctx, req)
if err != nil {
return blockRef, blockNumber, err
}

// Only rewrite moving tip-bound tags. Numeric refs, block-hash refs, "*",
// and slower-moving tags like "earliest" are already correct.
if blockRef != "latest" && blockRef != "finalized" && blockRef != "safe" {
return blockRef, blockNumber, nil
}

// WRITE path: prefer the response's own block number, which is the
// definitive answer for what payload we're about to cache.
if resp != nil {
if _, respBN, rerr := ExtractBlockReferenceFromResponse(ctx, resp); rerr == nil && respBN > 0 {
hex, herr := common.NormalizeHex(respBN)
if herr == nil {
return hex, respBN, nil
}
}
}

// READ path (and WRITE fallback): consult the network's aggregated view
// of the tag's current value. Guarded against panics because this helper
// is purely an optimization — if the network state isn't reachable for
// any reason (partially-constructed Network in a test, nil upstream
// registry, transient initialization race), we fall back to the tag-
// literal blockRef and retain the previous behaviour rather than
// aborting a live cache operation.
net := req.Network()
if net == nil {
return blockRef, blockNumber, nil
}
var num int64
func() {
defer func() {
if r := recover(); r != nil {
num = 0
}
}()
switch blockRef {
case "latest":
num = net.EvmHighestLatestBlockNumber(ctx)
case "finalized", "safe":
num = net.EvmHighestFinalizedBlockNumber(ctx)
}
}()
if num > 0 {
hex, herr := common.NormalizeHex(num)
if herr == nil {
return hex, num, nil
}
}

return blockRef, blockNumber, nil
}

func ExtractBlockReferenceFromResponse(ctx context.Context, r *common.NormalizedResponse) (string, int64, error) {
ctx, span := common.StartDetailSpan(ctx, "Evm.ExtractBlockReferenceFromResponse")
defer span.End()
Expand Down
82 changes: 82 additions & 0 deletions architecture/evm/block_ref_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -449,3 +449,85 @@ func TestExtractBlockReference(t *testing.T) {
})
}
}

// TestResolveCacheBlockRef covers the cache-specific helper that rewrites
// moving-tag blockRefs ("latest"/"finalized"/"safe") to a concrete block
// number so each tip advance produces a distinct cache key. The previous
// behaviour pinned all "latest" responses under the literal "latest" ref,
// causing stale cache hits for up to TTL after a tip advance.
func TestResolveCacheBlockRef(t *testing.T) {
ctx := context.Background()

t.Run("numeric ref passes through unchanged (no rewrite for non-tag)", func(t *testing.T) {
rpcReq := &common.JsonRpcRequest{
Method: "eth_getBlockByNumber",
Params: []interface{}{"0x1234", false},
}
nrq := common.NewNormalizedRequestFromJsonRpcRequest(rpcReq)

ref, num, err := ResolveCacheBlockRef(ctx, nrq, nil)
assert.NoError(t, err)
// ExtractBlockReferenceFromRequest normalizes a numeric request ref
// to its decimal string form; the helper forwards whatever that
// returns for non-tag refs.
assert.Equal(t, "4660", ref)
assert.Equal(t, int64(0x1234), num)
})

t.Run("latest tag + response with block number rewrites ref to response hex", func(t *testing.T) {
rpcReq := &common.JsonRpcRequest{
Method: "eth_getBlockByNumber",
Params: []interface{}{"latest", false},
}
nrq := common.NewNormalizedRequestFromJsonRpcRequest(rpcReq)
rpcResp := common.MustNewJsonRpcResponseFromBytes(nil, []byte(`{"number":"0xabcdef","hash":"0x1","parentHash":"0x0"}`), nil)
nrs := common.NewNormalizedResponse().WithJsonRpcResponse(rpcResp).WithRequest(nrq)
nrq.SetLastValidResponse(ctx, nrs)

ref, num, err := ResolveCacheBlockRef(ctx, nrq, nrs)
assert.NoError(t, err)
assert.Equal(t, "0xabcdef", ref, "write path must key by response's actual block number, not 'latest'")
assert.Equal(t, int64(0xabcdef), num)
})

t.Run("latest tag with no network and no response falls back to tag literal", func(t *testing.T) {
rpcReq := &common.JsonRpcRequest{
Method: "eth_getBlockByNumber",
Params: []interface{}{"latest", false},
}
nrq := common.NewNormalizedRequestFromJsonRpcRequest(rpcReq)

ref, num, err := ResolveCacheBlockRef(ctx, nrq, nil)
assert.NoError(t, err)
assert.Equal(t, "latest", ref, "with no response and no network, helper must fall back to tag so caller can decide to skip caching")
assert.Equal(t, int64(0), num)
})

t.Run("finalized tag rewrite on write path uses response block number", func(t *testing.T) {
rpcReq := &common.JsonRpcRequest{
Method: "eth_getBlockByNumber",
Params: []interface{}{"finalized", false},
}
nrq := common.NewNormalizedRequestFromJsonRpcRequest(rpcReq)
rpcResp := common.MustNewJsonRpcResponseFromBytes(nil, []byte(`{"number":"0x100","hash":"0x1","parentHash":"0x0"}`), nil)
nrs := common.NewNormalizedResponse().WithJsonRpcResponse(rpcResp).WithRequest(nrq)
nrq.SetLastValidResponse(ctx, nrs)

ref, num, err := ResolveCacheBlockRef(ctx, nrq, nrs)
assert.NoError(t, err)
assert.Equal(t, "0x100", ref)
assert.Equal(t, int64(0x100), num)
})

t.Run("earliest tag not rewritten (not tip-bound, existing semantics preserved)", func(t *testing.T) {
rpcReq := &common.JsonRpcRequest{
Method: "eth_getBlockByNumber",
Params: []interface{}{"earliest", false},
}
nrq := common.NewNormalizedRequestFromJsonRpcRequest(rpcReq)

ref, _, err := ResolveCacheBlockRef(ctx, nrq, nil)
assert.NoError(t, err)
assert.Equal(t, "earliest", ref, "earliest does not move with the tip; must not be rewritten")
})
}
4 changes: 3 additions & 1 deletion architecture/evm/error_normalizer.go
Original file line number Diff line number Diff line change
Expand Up @@ -309,7 +309,9 @@ func ExtractJsonRpcError(r *http.Response, nr *common.NormalizedResponse, jr *co
strings.Contains(ml, "already in the mempool") ||
strings.Contains(ml, "transaction already exists") ||
strings.Contains(ml, "already have transaction") ||
strings.Contains(ml, "already exists in mempool") {
strings.Contains(ml, "already exists in mempool") ||
strings.Contains(ml, "tx_replay_attack") ||
strings.Contains(ml, "replay attack") {
// These indicate the exact same transaction is already known - idempotent success case
return common.NewErrEndpointNonceException(
common.NewErrJsonRpcExceptionInternal(
Expand Down
52 changes: 52 additions & 0 deletions architecture/evm/error_normalizer_test.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package evm

import (
"errors"
"net/http"
"testing"

Expand Down Expand Up @@ -55,3 +56,54 @@ func TestExtractJsonRpcError_RequestTooLargeNormalization(t *testing.T) {
})
}
}

// TestExtractJsonRpcError_ReplayAttackIdempotency verifies that a re-submission
// of an already-accepted transaction rejected with a "replay attack" error is
// normalized to ErrEndpointNonceException with reason "already known", so
// eth_sendRawTransaction idempotency handling can convert it to success.
func TestExtractJsonRpcError_ReplayAttackIdempotency(t *testing.T) {
t.Parallel()

cases := []struct {
name string
message string
}{
{
name: "uppercase errmsg token",
message: "errcode: 113, errmsg: TX_REPLAY_ATTACK",
},
{
name: "spaced phrasing",
message: "transaction rejected: replay attack detected",
},
}

for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
t.Parallel()

r := &http.Response{StatusCode: 200, Header: http.Header{}}
jrErr := common.NewErrJsonRpcExceptionExternal(
int(common.JsonRpcErrorServerSideException),
tc.message,
"",
)
jr := common.MustNewJsonRpcResponse(1, nil, jrErr)

err := ExtractJsonRpcError(r, nil, jr, nil)
if err == nil {
t.Fatalf("expected error, got nil")
}
if !common.HasErrorCode(err, common.ErrCodeEndpointNonceException) {
t.Fatalf("expected ErrEndpointNonceException, got %T: %v", err, err)
}
var ne *common.ErrEndpointNonceException
if !errors.As(err, &ne) {
t.Fatalf("expected *common.ErrEndpointNonceException in chain, got %T", err)
}
if got := ne.Details["nonceExceptionReason"]; got != string(common.NonceExceptionReasonAlreadyKnown) {
t.Fatalf("expected reason %q, got %v", common.NonceExceptionReasonAlreadyKnown, got)
}
})
}
}
Loading
Loading