Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
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
1 change: 1 addition & 0 deletions Dockerfile
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,7 @@ COPY --from=builder /home/node/app/$location/package.json /home/node/app/$locati
# Copy Go binary and alias config
COPY --from=go-builder /build/streams-adapter /usr/local/bin/streams-adapter
COPY --from=builder /home/node/app/packages/streams-adapter/endpoint_aliases.json /home/node/app/endpoint_aliases.json
COPY --from=builder /home/node/app/packages/streams-adapter/adapter_includes.json /home/node/app/adapter_includes.json
COPY --from=builder /home/node/app/start-supervisor.sh /usr/local/bin/start-supervisor.sh

# Make scripts executable
Expand Down
57 changes: 56 additions & 1 deletion packages/scripts/src/generate-endpoint-aliases/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import path from 'path'
import { getWorkspaceAdapters } from '../workspace'

const OUTPUT_PATH = 'packages/streams-adapter/endpoint_aliases.json'
const INCLUDES_OUTPUT_PATH = 'packages/streams-adapter/adapter_includes.json'

/**
* Adapter types that are served by the streams adapter
Expand All @@ -23,7 +24,18 @@ interface EndpointConfig {
}

interface AllAdaptersConfig {
adapters: Record<string, { defaultEndpoint?: string; endpoints?: Record<string, EndpointConfig> }>
adapters: Record<
string,
{
defaultEndpoint?: string
endpoints?: Record<string, EndpointConfig>
includes?: Record<string, Record<string, { inverse: boolean }>>
}
>
}

interface AdapterIncludesConfig {
adapters: Record<string, Record<string, Record<string, { inverse: boolean }>>>
}

interface LoadResult {
Expand All @@ -48,6 +60,33 @@ async function loadAdapter(adapterPath: string): Promise<LoadResult> {
}
}

function extractIncludes(
adapter: Adapter,
): Record<string, Record<string, { inverse: boolean }>> | undefined {
const priceAdapter = adapter as Adapter & {
includesMap?: Record<string, Record<string, { inverse: boolean }>>
}
if (!priceAdapter.includesMap) {
return undefined
}

const includes: Record<string, Record<string, { inverse: boolean }>> = {}
for (const [from, toMap] of Object.entries(priceAdapter.includesMap)) {
if (!toMap || Object.keys(toMap).length === 0) {
continue
}
includes[from] = {}
for (const [to, details] of Object.entries(toMap)) {
if (!details) {
continue
}
includes[from][to] = { inverse: !!details.inverse }
}
}

return Object.keys(includes).length > 0 ? includes : undefined
}

function extractEndpoints(adapter: Adapter): Record<string, EndpointConfig> | undefined {
const endpoints: Record<string, EndpointConfig> = {}

Expand Down Expand Up @@ -91,6 +130,7 @@ async function main(): Promise<void> {
result.adapters[adapterKey] = {
defaultEndpoint: adapter.defaultEndpoint ?? undefined,
endpoints: extractEndpoints(adapter),
includes: extractIncludes(adapter),
}
} else {
skipped.push({ name: meta.descopedName, reason: skipReason || 'unknown' })
Expand All @@ -111,6 +151,21 @@ async function main(): Promise<void> {
fs.writeFileSync(outPath, JSON.stringify(result, null, 2), 'utf-8')
console.log(`Written ${Object.keys(result.adapters).length} EAv3 adapters to ${OUTPUT_PATH}`)

const includesResult: AdapterIncludesConfig = { adapters: {} }
for (const [adapterKey, adapterCfg] of Object.entries(result.adapters)) {
if (adapterCfg.includes) {
includesResult.adapters[adapterKey] = adapterCfg.includes
}
}

const includesOutPath = path.resolve(process.cwd(), INCLUDES_OUTPUT_PATH)
fs.writeFileSync(includesOutPath, JSON.stringify(includesResult, null, 2), 'utf-8')
console.log(
`Written ${
Object.keys(includesResult.adapters).length
} EAv3 adapters with includes to ${includesOutPath}`,
)

if (skipped.length > 0) {
console.log(`\nSkipped ${skipped.length} EAv3 adapters:`)
for (const { name, reason } of skipped) {
Expand Down
28 changes: 20 additions & 8 deletions packages/streams-adapter/cache/cache.go
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ import (
"github.com/prometheus/client_golang/prometheus/promauto"

types "streams-adapter/common"
"streams-adapter/includes"
)

var cacheDataGetCount = promauto.NewCounter(
Expand Down Expand Up @@ -62,13 +63,23 @@ type Cache struct {
items map[string]*types.CacheItem // rawKey → item
byTransformedKey map[string]map[string]struct{} // transformedKey → rawKeys (secondary index)
pendingObs map[string]*pendingObservation // transformedKey → buffered observation (pre-mapping race)
includes *includes.Index // adapter includes index for inverse flag lookup
ttl time.Duration
cleanupInterval time.Duration
ctx context.Context
cancel context.CancelFunc
stopOnce sync.Once
}

// SetIncludesIndex sets the adapter includes index used to determine the
// inverse flag from the original requested pair. When nil or the pair is not
// present, the cache defaults to not inverting.
func (c *Cache) SetIncludesIndex(idx *includes.Index) {
c.mu.Lock()
defer c.mu.Unlock()
c.includes = idx
}

// New creates a new cache instance
func New(cfg Config) *Cache {
ctx, cancel := context.WithCancel(context.Background())
Expand Down Expand Up @@ -141,7 +152,7 @@ func (c *Cache) SetTransformedKey(rawKey, transformedKey string) {
c.removeTransformedKeyMapping(item.TransformedKey, rawKey)
}
item.TransformedKey = transformedKey
item.RequiresInverse = requiresInverse(item.OriginalRequestData, transformedKey)
item.RequiresInverse = c.requiresInverse(item.OriginalRequestData)
item.Status = types.StatusLearned
item.Timestamp = time.Now()
c.addTransformedKeyMapping(transformedKey, rawKey)
Expand Down Expand Up @@ -262,7 +273,7 @@ func (c *Cache) RawKeysByTransformed(transformedKey string) ([]string, bool) {
return result, true
}

func requiresInverse(originalRequestData map[string]interface{}, transformedKey string) bool {
func (c *Cache) requiresInverse(originalRequestData map[string]interface{}) bool {
if originalRequestData == nil {
return false
}
Expand All @@ -273,14 +284,15 @@ func requiresInverse(originalRequestData map[string]interface{}, transformedKey
return false
}

transformedParams := parseCacheKey(transformedKey)
transformedBase := strings.ToUpper(transformedParams["base"])
transformedQuote := strings.ToUpper(transformedParams["quote"])
if transformedBase == "" || transformedQuote == "" {
return false
// The adapter_includes.json generated from the JS adapter's includes.json is
// the only source of truth for whether an observation must be inverted.
if c.includes != nil {
if inc, ok := c.includes.Lookup(originalBase, originalQuote); ok {
return inc.Inverse
}
}

return originalBase == transformedQuote && originalQuote == transformedBase
return false
}

func getPairValue(data map[string]interface{}, names ...string) string {
Expand Down
94 changes: 94 additions & 0 deletions packages/streams-adapter/cache/cache_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import (

types "streams-adapter/common"
helpers "streams-adapter/helpers"
"streams-adapter/includes"

"github.com/goccy/go-json"
"github.com/stretchr/testify/assert"
Expand Down Expand Up @@ -266,6 +267,99 @@ func TestCache_SetObservation_FansOutToSameTransformedKey(t *testing.T) {
}
}

func TestCache_SetTransformedKey_RequiresInverse_FromIncludes(t *testing.T) {
idx := includes.NewIndex(includes.AdapterIncludes{
"XAU": {"USD": {Inverse: false}},
"TRY": {"USD": {Inverse: true}},
})

cases := []struct {
name string
original map[string]interface{}
transformed string
wantInverse bool
}{
{
name: "metals swapped pair with inverse=false",
original: map[string]interface{}{"base": "XAU", "quote": "USD"},
transformed: "base=usd:endpoint=forex:quote=xau",
wantInverse: false,
},
{
name: "fiat swapped pair with inverse=true",
original: map[string]interface{}{"base": "TRY", "quote": "USD"},
transformed: "base=usd:endpoint=forex:quote=try",
wantInverse: true,
},
{
name: "direct pair not in includes defaults to false",
original: map[string]interface{}{"base": "USD", "quote": "TRY"},
transformed: "base=usd:endpoint=forex:quote=try",
wantInverse: false,
},
}

for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
c := New(Config{TTL: time.Minute, CleanupInterval: time.Hour})
defer c.Stop()
c.SetIncludesIndex(idx)

rawKey, err := helpers.CalculateCacheKey(types.RequestParams{
"endpoint": "forex", "base": tc.original["base"].(string), "quote": tc.original["quote"].(string),
})
require.NoError(t, err)

c.SetNew(rawKey, tc.original, [32]byte{})
c.SetTransformedKey(rawKey, tc.transformed)

item := c.Get(rawKey)
require.NotNil(t, item)
require.Equal(t, tc.wantInverse, item.RequiresInverse)
})
}
}

func TestCache_SetTransformedKey_RequiresInverse_DefaultFalseWhenNotInIncludes(t *testing.T) {
idx := includes.NewIndex(includes.AdapterIncludes{
"XAU": {"USD": {Inverse: false}},
})
c := New(Config{TTL: time.Minute, CleanupInterval: time.Hour})
defer c.Stop()
c.SetIncludesIndex(idx)

rawKey, err := helpers.CalculateCacheKey(types.RequestParams{
"endpoint": "forex", "base": "EUR", "quote": "USD",
})
require.NoError(t, err)
transformed := "base=usd:endpoint=forex:quote=eur"

c.SetNew(rawKey, map[string]interface{}{"base": "EUR", "quote": "USD"}, [32]byte{})
c.SetTransformedKey(rawKey, transformed)

item := c.Get(rawKey)
require.NotNil(t, item)
require.False(t, item.RequiresInverse, "pair not in includes must not be inverted")
}

func TestCache_SetTransformedKey_RequiresInverse_DefaultFalseWithoutIndex(t *testing.T) {
c := New(Config{TTL: time.Minute, CleanupInterval: time.Hour})
defer c.Stop()

rawKey, err := helpers.CalculateCacheKey(types.RequestParams{
"endpoint": "forex", "base": "EUR", "quote": "USD",
})
require.NoError(t, err)
transformed := "base=usd:endpoint=forex:quote=eur"

c.SetNew(rawKey, map[string]interface{}{"base": "EUR", "quote": "USD"}, [32]byte{})
c.SetTransformedKey(rawKey, transformed)

item := c.Get(rawKey)
require.NotNil(t, item)
require.False(t, item.RequiresInverse, "pair without an includes index must not be inverted")
}

func TestCache_SetTransformedKey_UsesExistingObservationForSharedTransformedKey(t *testing.T) {
c := New(Config{TTL: time.Minute, CleanupInterval: time.Hour})
defer c.Stop()
Expand Down
3 changes: 3 additions & 0 deletions packages/streams-adapter/config/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,9 @@ type Config struct {
// Other configuration
LogLevel string
AdapterName string

// Version is populated at runtime by the JS adapter health endpoint.
Version string
}

// Load reads configuration from environment variables
Expand Down
82 changes: 82 additions & 0 deletions packages/streams-adapter/helpers/inversion.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,82 @@
package helpers

import (
"encoding/json"
"fmt"
"strconv"

types "streams-adapter/common"
)

// InvertObservation returns a copy of obs with its numeric result(s) replaced
// by their reciprocal. Used for requests whose original pair is the inverse
// of the transformed key the provider actually publishes.
func InvertObservation(obs *types.Observation) (*types.Observation, error) {
inverted := *obs

data, err := invertResultInObject(obs.Data)
if err != nil {
return nil, err
}
inverted.Data = data

if len(obs.Result) > 0 {
result, err := invertRawNumber(obs.Result)
if err != nil {
return nil, err
}
inverted.Result = result
}

return &inverted, nil
}

func invertResultInObject(raw json.RawMessage) (json.RawMessage, error) {
var data map[string]interface{}
if err := json.Unmarshal(raw, &data); err != nil {
return nil, fmt.Errorf("unable to invert observation result: %w", err)
}

result, ok := data["result"]
if !ok {
return nil, fmt.Errorf("unable to invert observation result: missing result")
}
num, err := numberFromInterface(result)
if err != nil {
return nil, err
}
if num == 0 {
return nil, fmt.Errorf("unable to invert observation result: result is zero")
}

data["result"] = 1 / num
return json.Marshal(data)
}

func invertRawNumber(raw json.RawMessage) (json.RawMessage, error) {
var num float64
if err := json.Unmarshal(raw, &num); err != nil {
return nil, fmt.Errorf("unable to invert top-level result: %w", err)
}
if num == 0 {
return nil, fmt.Errorf("unable to invert top-level result: result is zero")
}
return json.Marshal(1 / num)
}

func numberFromInterface(value interface{}) (float64, error) {
switch v := value.(type) {
case float64:
return v, nil
case json.Number:
return v.Float64()
case string:
num, err := strconv.ParseFloat(v, 64)
if err != nil {
return 0, fmt.Errorf("unable to invert observation result: result is not numeric")
}
return num, nil
default:
return 0, fmt.Errorf("unable to invert observation result: result is not numeric")
}
}
Loading
Loading