Skip to content

Caching

Mycel provides two flow-level caching mechanisms: cache for avoiding repeated reads by storing results, and dedupe (since v2.1.0) for dropping no-op writes whose persisted projection is byte-identical to the last one processed for the same key.

Cache Setup

First, define a cache connector:

# Redis (recommended for production and multi-instance)
connector "redis_cache" {
  type    = "cache"
  driver  = "redis"
  url     = env("REDIS_URL", "redis://localhost:6379")
  prefix  = "myapp:"
}

# In-memory (development or single-instance)
connector "local_cache" {
  type      = "cache"
  driver    = "memory"
  max_items = 10000
  eviction  = "lru"
}

Inline Cache Block

Add a cache block directly in a flow:

flow "get_product" {
  from {
    connector = "api"
    operation = "GET /products/:id"
  }

  cache {
    storage = "redis_cache"
    ttl     = "5m"
    key     = "product:${input.id}"
  }

  to {
    connector = "db"
    target    = "products WHERE id = :id"
  }
}

When a request comes in: 1. Mycel computes the cache key 2. If the key exists in the cache, return the cached value immediately (no to block executes) 3. If not, execute the to block and store the result in the cache

Cache Attributes

Attribute Type Required Description
storage string yes Cache connector name
ttl string no Time-to-live: "5m", "1h", "24h"
key string no Key template — ${...} is substituted, the rest is literal (default: auto-generated from the request)
key_from string no CEL expression yielding the key, evaluated against input.* before the lookup — for a list or map input that has to be sorted, joined or hashed first. Mutually exclusive with key. See Deriving the key from a list or map
invalidate_on list no Event patterns that invalidate this cache entry
use string no Reference a named cache (use = "cache.<name>"); its storage, ttl, prefix and encoding come with it
encoding list no Wire format for entries, applied in order on the way out and reversed on the way in. Default ["json"]. See Sharing a namespace

Cache Key Templates

The cache key must uniquely identify the request. It is a template, not a CEL expression: each ${...} is replaced by the value it names, and everything outside is the key as written.

# Simple ID-based key
key = "product:${input.id}"

# Multiple parameters
key = "users:${input.id}:orders:${input.status}"

# Per-caller cache, keyed by a request header
key = "user_data:${input.headers.x-user-id}"

Not CEL — the lock, dedupe and coordinate keys are, this one is not

key = "'product:' + input.id" is the form the sync primitives take, and it is accepted here without complaint: the key is used verbatim, quotes and + included, so every request shares one cache entry and gets back whichever record was cached first, for the life of the TTL. Nothing fails — the symptom is users seeing the wrong record. Since 3.6.2 mycel validate refuses a key whose text outside ${...} carries quotes, a +, or an input. reference, and shows the template form to write instead. When the key genuinely needs an expression, that is what key_from is for.

Deriving the key from a list or map — key_from

A template substitutes scalars. When what identifies the request is a list or a map — a faceted listing's filter: [{attribute_code, value}, ...] — there is nothing a template can say: ${input.filter} renders Go's own syntax into the key ([map[attribute_code:room value:Bath]], with a warning in the log), and the canonical form the key needs cannot be expressed before the lookup runs.

key_from takes a CEL expression instead, evaluated against input.* before the lookup, with the same scope and rules as keys_from on the invalidation side:

flow "gallery" {
  from {
    connector = "api"
    operation = "Query.items"
  }

  cache {
    storage  = "redis_cache"
    ttl      = "1h"
    key_from = <<-CEL
      'gallery:' + (input.store ?? 'default') + ':' +
      hash_sha256(join(as_list(input.filter).map(f, f.attribute_code + '=' + f.value), '|'))
    CEL
  }

  step "rows" {
    connector = "db"
    query     = "SELECT * FROM items"
  }
}

The same filter set is the same key, and two different sets cannot collide. The rules:

  • key and key_from are mutually exclusive; the parser refuses a block with both, and mycel validate refuses a key_from that carries ${...}, since that is the template form in the wrong attribute.
  • The expression must yield a non-empty string. A list, a number or an empty string fails the request rather than caching under an empty or literal key, and the error names the flow and the attribute. mycel validate does not evaluate CEL, so a wrong expression shows up on the first request, and on every request after it.
  • A named cache's prefix goes in front of a derived key the same as a template one.
  • An aspect's cache {} block takes key_from too, with the same rules; there key is no longer required when key_from is given.

join, as_list and hash_sha256 are the functions this usually needs; the order of the list is the order of the key, so sort it upstream if two orderings of the same set should hit one entry.

Named Caches

Define a named cache block to reuse cache configuration across multiple flows:

cache "products" {
  storage       = "redis_cache"
  ttl           = "10m"
  prefix        = "products"
  invalidate_on = ["product.updated", "product.deleted"]
}

Reference it in a flow:

flow "get_product" {
  from {
    connector = "api"
    operation = "GET /products/:id"
  }
  cache = cache.products
  to {
    connector = "db"
    target    = "products WHERE id = :id"
  }
}

Named Cache Attributes

Attribute Type Description
storage string Cache connector name
ttl string Default TTL for entries
prefix string Key prefix for namespacing
invalidate_on list Event patterns that trigger invalidation
encoding list Wire format for entries in this cache, inherited by any flow that references it

Cache Invalidation

invalidate_on (automatic)

Invalidate cache entries when specific events happen. This uses event pattern matching:

flow "get_user" {
  from {
    connector = "api"
    operation = "GET /users/:id"
  }

  cache {
    storage       = "redis_cache"
    ttl           = "15m"
    key           = "user:${input.id}"
    invalidate_on = ["user.updated:${input.id}", "user.deleted:${input.id}"]
  }

  to {
    connector = "db"
    target    = "users WHERE id = :id"
  }
}

after block (explicit, per-mutation flow)

Explicitly invalidate keys after a write operation:

flow "update_product" {
  from {
    connector = "api"
    operation = "PUT /products/:id"
  }
  to {
    connector = "db"
    target    = "UPDATE products"
  }

  after {
    invalidate {
      storage  = "redis_cache"
      keys     = ["product:${input.id}"]
      patterns = ["products:list:*"]
    }
  }
}

keys invalidates exact keys. patterns invalidates all matching keys (glob-style). Both are templates: ${input.id} is replaced when the flow runs.

A key set whose size depends on the data

keys is one key out per template in. The values in a key vary with the message; the number of keys is fixed when the configuration is parsed. When the set is whatever a query returned — every store view a product appears in, every rewrite path it has had, every variant of a parent — that number is a function of the data, and there is nothing to write.

keys_from and patterns_from take a CEL expression yielding a list of strings:

flow "republish_product" {
  from {
    connector = "api"
    operation = "PUT /products/:id/republish"
  }

  step "affected" {
    connector = "db"
    query     = "SELECT store_code FROM product_stores WHERE product_id = :id"
    params    = { id = "input.id" }
  }

  to {
    connector = "db"
    target    = "products"
  }

  after {
    invalidate {
      storage   = "redis_cache"
      keys      = ["products:item"]
      keys_from = "step.affected.map(r, 'product:' + input.id + ':' + r.store_code)"
    }
  }
}

input.*, output.* and step.* are in scope, because the list almost always comes from a query the flow just ran. The result is unioned with the static list and deduplicated, so a fixed key and a computed set can be named together.

A wildcard is not a substitute

Not when the members diverge. Rewrite paths drift from the URL key through redirects and history, so a prefix broad enough to catch them all also deletes entries for unrelated products, and one narrow enough to be safe misses exactly the paths that most need dropping — over-invalidating and under-invalidating at the same time.

Aiming a ${...} template at a list does not fan out either: it renders Go syntax into the key (url-rewrite-[a b c]), deletes that, and reports success. It now warns and points here.

When an invalidation does not happen

A flow whose invalidation failed has already done and committed its own work, so it answers 200 either way — and the cache is now serving what that write made stale, with nothing to correct it. The two ways it can fail want different answers:

What failed Response
The cache could not be reached (keys, patterns) Logged at warn, counted as mycel_cache_invalidate_errors_total. The request still succeeds — the write is committed, and failing it afterwards would be wrong
keys_from / patterns_from could not be evaluated, or did not yield a list of strings Logged, counted, and the request fails. This is not transient: it will fail identically on every message for as long as the flow is deployed, invalidating nothing while reporting success, and mycel validate does not evaluate CEL so it is not caught beforehand either
WARN cache invalidation did not happen
     flow=invalidate_products cache=redis_cache attr=keys_from fatal=true
     error="invalidate keys_from: CEL eval error: no such key: nope"

This matters most for a flow whose only job is invalidation — an endpoint a consumer calls after writing elsewhere, with steps and an after block and no to. There is nothing else to observe, so a silent no-op would answer 200 forever and the only symptom would be stale reads somewhere else entirely, hours later. examples/cache shows that shape as Pattern 7.

Numbers read back from the cache

An entry is stored as JSON and the digits are kept as written: an integral number reads back as an int64 when it fits one, a fraction as a float64. Before 3.7.0 every number came back as a float64, so an integer past 2^53 — a snowflake id, a 64-bit hash, cents in a large ledger — was served from the cache with its low digits rounded away while the first, uncached answer had them right.

Sharing a namespace with another service

A flow's cache entries are ["json"] unless the block says otherwise, and that is the right answer while Mycel owns the namespace. It stops being the right answer during a migration — which is exactly when a cache is most likely to be shared, because the service being replaced is still up and still reading and writing the same keys.

Getting it wrong is not incompatibility, it is mutual destruction:

  1. Mycel reads a key the other service wrote and cannot decode it.
  2. That reads as a miss, so the flow does the work.
  3. The flow then writes its format over that key.
  4. The other service reads the same key next and its own decode throws.

They take turns destroying each other's entries, and the only visible symptom is a cache that never seems to hit.

encoding declares the format. The codecs listed apply left to right on the way out and reversed on the way in:

cache {
  storage  = "redis_cache"
  ttl      = "5m"
  key      = "product:${input.id}"
  # Reads and writes what a service storing gzip(base64(JSON.stringify(v))) does
  encoding = ["json", "base64", "gzip"]
}
Codec Position What it does
json first, required The value ↔ bytes
base64 after Bytes ↔ bytes, standard alphabet with padding. Tolerant on the way in of missing padding and wrapped lines
gzip after Bytes ↔ bytes

A chain that could never be applied — one that does not start with json, or names a codec that does not exist — fails when the configuration is read, not on the first cache write.

A named cache can carry it, so flows sharing a namespace share its format without restating it; a flow declaring its own wins.

Knowing when the format is wrong

A found entry that cannot be decoded is not a hit. It gets its own counter — mycel_cache_decode_errors_total — because it is neither a hit (the flow is about to do the work) nor a miss the cache could fix by being warmer, and it is logged at warn naming the flow, the cache and the key:

WARN cache entry could not be decoded; treating as a miss and doing the work
     flow=get_product cache=redis_cache key=product:42 bytes=180

That line is the only signal that the cache holds something this flow cannot use. Note that a hit rate computed as hits / (hits + misses) will not show this — see Observability.

Deduplication

Since v2.1.0 the dedupe block is content-based and runs in two phases. Phase A (after transform, before to) computes a canonical fingerprint over the projection the operator declares and compares it byte-for-byte to the stored fingerprint for the same key; on match the message is dropped according to on_duplicate without invoking to. Phase B (after to succeeds) stores the new fingerprint, so a failed-then-retried message will not self-discard.

The primitive self-locks per key (in-process via the memory-backed SyncManager) so two workers cannot both pass Phase A with identical fingerprints and double-call the downstream. For cross-process serialization across multiple Mycel pods, compose with an outer lock {} block on the same resource key.

The typical use case is an MQ consumer where the upstream re-sends "update" messages even when nothing relevant changed: every redelivery hits a slow downstream and the queue accumulates. With dedupe, only messages whose persisted projection actually differs reach the downstream.

connector "fp_cache" {
  type   = "cache"
  driver = "redis"   # or "memory" for tests / single-pod
}

flow "process_payment" {
  from {
    connector = "rabbit"
    operation = "payments"
  }

  transform {
    payment_id = "input.payment_id"
    account_id = "input.account_id"
    amount     = "input.amount"
  }

  dedupe {
    cache        = "fp_cache"
    key          = "'payment:' + input.payment_id"
    ttl          = "24h"
    on_duplicate = "ack"
    fingerprint {
      payment_id = "output.payment_id"
      account_id = "output.account_id"
      amount     = "output.amount"
    }
  }

  to {
    connector = "db"
    target    = "payments"
  }
}

Dedupe Attributes

Attribute Type Required Default Description
cache string yes — Name of a connector { type = "cache" }. The connector pool is initialized once at startup; the hot path does not pay a registry lookup per message
key string yes — CEL expression for the per-resource fingerprint key (evaluated against input.*)
fingerprint {} block yes* — Named CEL expressions whose values form the projection. Both input.* and output.* (transform result) are in scope. Must list every persisted field — omitting one would silently drop real changes. *Required unless the block declares facet blocks instead
facet "<name>" {} block no — An independently-tracked part of the projection, holding its own fingerprint {}. Each is stored and committed on its own; a to naming it runs only when it changed, and the message is dropped only when no facet did. Cannot be combined with a bare fingerprint {}. See Splitting a message into parts
ttl string no — How long to keep stored fingerprints. Supports "30d" and "2w" plus stdlib units (s/m/h); malformed values fail the parse
on_duplicate string no "ack" Behavior on fingerprint match: "ack", "reject", "requeue". Matches the sequence_guard vocabulary so MQ consumers handle it uniformly
compare_when string no — CEL predicate gating Phase A only. False: the stored fingerprint is not consulted, so the message cannot be dropped; Phase B still commits after a successful write. input.* and output.* in scope. See When the record can vanish

When the record can vanish

A stored fingerprint says the content was written once. It does not say it is still written. Nothing in dedupe observes the downstream record being removed by a path the flow never sees — a manual delete in an admin UI, a restore, a data fix — and there is usually no flow to clear the fingerprint when that happens. The re-send meant to repair the damage then matches a fingerprint describing a record that no longer exists, and is dropped.

compare_when gates the comparison, and only the comparison:

step "check_present" {
  connector = "db"
  query     = "SELECT CAST(COALESCE((SELECT 1 FROM products WHERE sku = :sku LIMIT 1), 0) AS SIGNED) AS row_exists"
  params    = { sku = "input.body.sku" }
  on_error  = "fail"
}

transform {
  row_exists = "int(step.check_present.row_exists)"
  name       = "input.body.name"
}

dedupe {
  cache        = "fp_cache"
  key          = "'sku_fp:' + input.body.sku"
  ttl          = "30d"
  compare_when = "output.row_exists == 1"
  fingerprint {
    name = "output.name"
  }
}
  • false → Phase A is skipped entirely: no GET, no comparison, no drop. Phase B still runs, so the write commits a fresh fingerprint and the next message can be suppressed. Gating both phases would leave the cache empty forever and the primitive permanently inert on the flow.
  • true or absent → unchanged behavior.
  • Fails open: a predicate that cannot be evaluated, or that does not return a boolean, logs a warning and processes the message. Same direction as the cache-error path — one extra downstream call is recoverable, a silently swallowed message is not.

The existence check goes in compare_when, not in fingerprint {}

Adding row_exists to the projection looks like it should work — record gone, projection differs, message reprocessed — and it does the opposite in both directions.

Phase A and Phase B share one projection: the fingerprint Phase B stores is the one Phase A computed, i.e. the pre-write reading. On a create, "does this record exist" is 0 by definition, so the stored value is 0 forever while every later message computes 1.

situation stored computed result
record exists, duplicate re-send 0 1 mismatch → duplicate reaches the destination
record deleted externally, re-send 0 0 match → dropped, which is the case it was added for

It breaks suppression exactly where suppression worked, and stays inert exactly where invalidation was needed.

Splitting a message into parts

A fingerprint answers one question — did anything change — and its only verb is to drop the message. That is right while every part of a message costs about the same to apply. It stops being right when it does not.

A product message carries its data and its images together. The data is a handful of columns; the images are URLs the downstream goes and fetches, which is minutes of work. With one fingerprint, a message whose name changed re-sends the images too, and anything waiting on that product waits for them to finish.

facet blocks split the projection into parts that are fingerprinted, stored and committed independently. A to naming a facet runs only when that facet changed; the message is dropped only when none did.

dedupe {
  cache        = "fp_cache"
  key          = "'item_fp:' + input.body.sku"
  ttl          = "30d"
  on_duplicate = "ack"              # applies when no facet changed

  facet "data" {
    fingerprint {
      sku   = "output.sku"
      name  = "output.name"
      price = "output.price"
    }
  }

  facet "assets" {
    fingerprint {
      sku        = "output.sku"
      main_image = "output.main_image"
      gallery    = "output.gallery"
    }
  }
}

to {
  facet     = "data"                # skipped when the data facet is unchanged
  connector = "catalog"
  target    = "catalog_items"
}

to {
  facet     = "assets"              # skipped when the assets facet is unchanged
  connector = "asset_jobs"          # a queue a second consumer reads
  target    = "item.assets.q"
}

What each kind of message now does:

The message changes data destination assets destination
Nothing — — (message dropped)
The name only writes —
An image only — publishes
Both writes publishes

Commit is per facet, and that is the point. A facet is committed only once every destination naming it has succeeded. If the catalogue write lands and the queue publish fails, the data facet is committed and the assets facet is not: the retry finds the data unchanged, does not write it a second time, and publishes the assets. Committing them together would have lost the images for the life of the entry — the same failure the two-phase commit exists to prevent, one level down.

What a facet's success means. It means its to succeeded, and nothing more. For a queue that is "the broker accepted the message", not "the other consumer finished the work" — there is no callback, and Mycel does not ask. The assets facet records that this set of images was handed over once. Past the handoff it is the second consumer's problem: its retries, its dead-letter queue, its own dedupe. Design the split so that is a reasonable place for the responsibility to change hands.

A few rules worth knowing before reaching for them:

  • A to without facet runs whenever the message is not dropped — which is what every flow without facets does, so nothing changes for them.
  • Each facet is stored under its own key, so a facet that has never been seen reads as changed. Adding a facet to a live flow therefore re-runs its destinations once per key: a backfill, not a bug, and worth planning the rollout around.
  • compare_when stays a single flow-level gate. When it is false nothing is compared and every facet runs — which is what a missing downstream record needs, since an assets-only message against a record that no longer exists would otherwise never re-create it.
  • A bare fingerprint {} and facet blocks in the same dedupe are refused: honouring both would mean deciding which one drops the message.
  • mycel validate refuses a to naming a facet nobody declared, and a facet no destination names. Both would otherwise be silent — the first skipped on every message, the second never committed and therefore permanently "changed", which quietly disables deduplication for the whole flow.

Mycel does not check that facets are independent

Facets are a statement by the author that these parts of the message can be applied separately. Two facets whose destinations write the same record will race, and nothing here will report it. Split by what the destinations actually touch — a domain, a table, a subsystem — not by what is convenient to name.

Pipeline order

The dedupe block runs after transform. The fingerprint expressions reference output.* (the transformed payload), so transform must run first. Earlier versions (≤ 2.0.0) had a key-based dedupe block that ran before transform; see CHANGELOG v2.1.0 for migration.

Array order-insensitivity

The canonical encoder sorts array elements before serialization, treating them as order-insensitive sets. This is appropriate for projections like "list of attribute values" or "set of website flags," but lossy for fields where order is semantically meaningful (e.g. a ranked list where position encodes priority).

For order-sensitive arrays, reshape them in transform before dedupe sees them — join with a delimiter into a single string:

transform {
  # Bad: ranked_tags as an array would lose order in the fingerprint.
  # Good: join into a string so order is part of the encoded value.
  ranked_tags = "input.ranked_tags.map(t, t).join(',')"
}

Caching vs Deduplication

Cache Dedupe
Purpose Avoid redundant downstream reads Drop no-op writes
Applies to Read flows Write flows (especially MQ consumers)
Cache miss Execute to, cache result Process normally; store fingerprint after to success
Cache hit Return cached value immediately Drop without invoking to
Compares Key only Canonical content fingerprint
Pipeline position Before to (read path) After transform, before to

Production Considerations

  • Use Redis for multi-instance deployments. In-memory cache is not shared across instances.
  • Set TTLs appropriate to your data freshness requirements. Stale cache is worse than no cache for critical data.
  • Use invalidate_on or the after block to invalidate caches on writes.
  • Monitor cache hit rates with the /metrics endpoint (Prometheus).
  • Use prefix or key expressions to prevent key collisions between services sharing a Redis instance.

Example: Read-Through Cache for Product Catalog

connector "redis_cache" {
  type    = "cache"
  driver  = "redis"
  url     = env("REDIS_URL", "redis://localhost:6379")
  prefix  = "catalog:"
}

# Cache product reads for 10 minutes
flow "get_product" {
  from {
    connector = "api"
    operation = "GET /products/:id"
  }

  cache {
    storage = "redis_cache"
    ttl     = "10m"
    key     = "product:${input.id}"
  }

  to {
    connector = "db"
    target    = "products WHERE id = :id"
  }
}

# Invalidate on update
flow "update_product" {
  from {
    connector = "api"
    operation = "PUT /products/:id"
  }
  to {
    connector = "db"
    target    = "UPDATE products"
  }

  after {
    invalidate {
      storage = "redis_cache"
      keys    = ["product:${input.id}"]
    }
  }
}

# Deduplicate no-op inventory updates by content
flow "handle_inventory_update" {
  from {
    connector = "rabbit"
    operation = "inventory.updated"
  }

  transform {
    product_id  = "input.product_id"
    stock_qty   = "input.stock_qty"
    reorder_at  = "input.reorder_at"
  }

  dedupe {
    cache        = "redis_cache"
    key          = "'inv_fp:' + input.product_id"
    ttl          = "1h"
    on_duplicate = "ack"
    fingerprint {
      product_id = "output.product_id"
      stock_qty  = "output.stock_qty"
      reorder_at = "output.reorder_at"
    }
  }

  to {
    connector = "db"
    target    = "UPDATE products"
  }
}

See Also