Skip to content

Feat/117 mandi price catalog crawler - #54

Open
kelvinprabhu wants to merge 67 commits into
developmentfrom
feat/117-Mandi-Price-Catalog-Crawler
Open

kelvinprabhu wants to merge 67 commits into
developmentfrom
feat/117-Mandi-Price-Catalog-Crawler

Conversation

@kelvinprabhu

@kelvinprabhu kelvinprabhu commented Sep 25, 2026 •

Copy link
Copy Markdown
Collaborator

Registry-Driven Catalogue Publishing

Summary

Publishing is now registry-driven. The registry decides which capabilities publish and which pipeline they use. The crawler discovers these bindings, runs pipelines on schedule, and publishes both crawled and pipeline-generated catalogues through a common sink to the provider adapter's /publish.

What Changed

Crawler

  • Added buildPublishSource / publishDiscoverer.Discover.
  • Discovers all ProviderSchema bindings with an active publish action and an embedded pipeline.
  • Added scheduled publish sweeps in publishpipelines.go.
  • Cron + PostgreSQL run log determine whether a pipeline is due.
  • Added overlap protection; one pipeline failure does not stop others.
  • publishPipelines: "true" enables publishing.
  • publishBindingKeys is now rejected.

Sink

  • BuildPushBody uses action: catalog/publish.
  • Adds schemaTypes to publishDirectives.
  • Publishes using MERGE instead of FULL.
  • Client.Push requires HTTP 200 + ACCEPTED on_publish verdict.
  • Added DiscoverySink.Publish implementing pipeline.Publisher.
  • Pipeline and crawled catalogues now share the same publishing path.

Pipeline

  • HTTP publishing was removed from the pipeline frame; RunOptions.Publisher handles it.
  • Removed catalogpublisher/catalogpublish; mapper helpers moved to pipeline/mapper.go.
  • upstream: is optional; pipelines without it skip credentials/token exchange.
  • Added const step for inline records.
  • Publish-address errors now reference the pipeline's own configuration.

Registry & Embedding

  • Added ProviderBindingLister / ProviderBindingKeys.
  • Registry failures are errors, not empty binding lists.
  • Publishing pipelines are centrally embedded under:
    pkg/plugin/implementation/*/cataloguepublish-*/
  • Replaces per-capability pipeline_files.go and the crawler collectors table.
  • Agmarknet uses CATALOG_PUBLISH_URL.

Breaking Changes

  • discoveryPushUrl now points to the provider adapter /publish, e.g.:

    discoveryPushUrl: "http://provider-adapter:9200/publish"
  • URLs ending in /push are rejected.

  • Crawled catalogues use MERGE; large catalogues may be split into multiple requests.

  • publishBindingKeys is removed; use publishPipelines: "true".

  • Every registry binding with a publish action runs on the node; there is no binding filter.

  • MANDI_PUBLISH_URL → CATALOG_PUBLISH_URL.

Configuration

plugins:
  registry:
    id: sunbirdRegistry
    config:
      url: http://sunbird-registry-service:8081/api/v1
      entity: Participant
      providerEntity: ProviderSchema

  crawler:
    id: catalogcrawler
    config:
      dbDsn: "postgres://crawler:crawler@crawler-db:5432/catalogcrawler?sslmode=disable"
      discoveryPushUrl: "http://provider-adapter:9200/publish"
      publishPipelines: "true"
      publishEnabled: "false"

A registry binding enables a pipeline with:

{
  "action": "publish",
  "mappings": "pkg/plugin/implementation/<Capability>/cataloguepublish-<source>/<pipeline>.yaml",
  "status": "active"
}

Testing

  • go build ./... and go vet pass.
  • All relevant tests pass; crawler tests pass with -race.
  • Sink tests cover ACCEPTED, PARTIAL, REJECTED, malformed/truncated responses, batching, HTTP errors, and schemaTypes.
  • Crawler tests cover discovery filtering, registry failures, overlap handling, retries, and URL validation.
  • All embedded pipelines validate against the schema.
  • Live Agmarknet test published 10 catalogues / 1,341 markets / 36 states in the first sweep; subsequent sweeps correctly skipped already-completed runs.

Adding a Publishing Capability

  1. Create:
    pkg/plugin/implementation/<Capability>/cataloguepublish-<source>/

  2. Add the pipeline YAML and mappings/.

  3. Run:

    go test ./pkg/plugin/implementation/
  4. Add the publish action to the ProviderSchema registry record.

  5. Rebuild.

The registry is now the single source of truth for publishing capability selection.

Comment thread pkg/plugin/implementation/catalogpublisher/pipeline/catalog.go Fixed
Comment thread pkg/plugin/implementation/catalogpublisher/pipeline/expr.go Fixed
@github-actions

github-actions Bot commented Sep 25, 2026 •

Copy link
Copy Markdown

🛡️ Trivy security scan (CRITICAL,HIGH,MEDIUM,LOW)

View full run

Go dependencies

No findings at CRITICAL,HIGH,MEDIUM,LOW.

Container image

No findings at CRITICAL,HIGH,MEDIUM,LOW.

@github-actions

github-actions Bot commented Sep 25, 2026 •

Copy link
Copy Markdown

📊 Test Coverage: ✅ Passed — 85% of changed lines covered, min 80%

@manjudr manjudr left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Review

Reviewed against the stated intent: enhance the crawler to support API-based capabilities + Sunbird registry compatibility, transformation logic in JSONata, pre-requisites injectable from outside the crawler package, so that the next capability (e.g. weather) is a config change rather than a Go change.

Checked out the branch, built it, and ran the PR-touched packages — go build ./... is clean and all tests pass (catalogpublisher, catalogpublisher/pipeline, catalogcrawler, sink, source, store, MandiPrice, implementation, sunbirdRegistry). Live tests are correctly gated behind MANDI_LIVE=1 / SUNBIRD_REGISTRY_URL, so go test ./... stays hermetic. Nice.

The architecture is right and the engineering quality is high — registry-as-source-of-truth, embedded pipelines, one sink shared by both the crawl and publish paths, cron + run-log scheduling. My concern is not the design; it's that a few things the YAML claims are not what the Go actually does, and the "next capability with no Go changes" promise doesn't hold yet.

Against the requirements

Requirement Status
API-based capability ✅ http.get / http.post, token exchange, scheduled, works live
Sunbird registry compatibility ✅ ProviderBindingLister; registry decides what publishes
Configured schedule interval ✅ cron + timezone in YAML, 5-min due-check, Postgres run log
Transformation logic in JSONata ⚠️ Partial
Pre-reqs outside crawler, easy to inject ❌ Outside, yes. Injectable, no

On JSONata: request/response shaping and catalogue rendering genuinely are JSONata (mappings/*.yaml) — that part is well done. But join / derive / dedupe / filter / chunk / exclude / order are a bespoke YAML mini-language executed by Go, and With (spec.go:189) is explicitly "the union of every field any step uses". Adding a primitive or even a new step param means editing spec.go + pipeline.v1.schema.json + steps.go. So "a folder and a registry entry, not Go" holds only for a capability shaped exactly like Agmarknet.


Blockers

1. Every upstream is forced through a token exchange — this blocks the weather case directly.
run.go:246 requires baseUrl + tokenUser + tokenSecret to be non-empty whenever an upstream: block exists; run.go:286 then calls client.Token(...) unconditionally; upstream.go:88 rejects any auth.kind other than tokenExchange. There is no none, apiKey, basic, or OAuth2 client-credentials. A weather API on a static Bearer key or an x-api-key header cannot be expressed — and dropping the upstream: block turns off HTTP steps entirely (hasUpstream gates the client construction). Suggest a small Authenticator interface selected on auth.kind, with none and apiKey as the first two implementations.

2. reauth is classified and then discarded.
The YAML declares reexchangeOn: [401, 403], but ReexchangeOn is read by no code anywhere. classifyReauth returns a plain error, which isn't ErrNoUpstreamData, so steps.go:510 buckets it as transportError → stateErrors++ → refuseWhen: collection.stateErrors > 0 refuses the entire day's publish. A token with no expiry field held across 36 sequential state calls is exactly the case this will bite. Either implement the re-exchange or delete the key.

3. The run log is keyed on capability, not pipeline.
run.go:224 writes RecordPipelineRun(ctx, capability, now) into a table whose PK is pipeline text. The cataloguepublish-<source> folder convention explicitly anticipates more than one source per capability — two would share one row and permanently starve each other. Key on the pipeline path.

4. No distributed lease.
publishSweep's inFlight is an in-process mutex, and the run mark is written only after a run that takes minutes. Two crawler replicas both read "not run yet", both execute the full pipeline, both publish. Worth claiming the row first with a conditional UPDATE ... WHERE last_run_at < $window RETURNING.


Should fix before merge

  • mandi-price-agmarket.yaml header (lines 12–22) is now false. It declares upstream, pipeline, discover and registry as INERT and points at collect.go — which this PR deletes. Those blocks are live now. This is the file the weather author will copy, so the stale header will propagate.
  • Inert spec fields — schema-validated, read by nothing: all three upstream.guards.* (Go hardcodes maxResponseBytes as a const in upstream.go), auth.token.reexchangeOn, and auth.token.name (the header path hardcodes Authorization: Bearer, the POST path hardcodes q.Set("token", ...)), plus publish.concurrency and publish.timeout. This is precisely the failure mode spec.go's own doc comment says it exists to prevent — "a rule written in the file can be believed to be in force while no code ever sees it."
  • The existing crawl path silently switched FULL → MERGE. Deletions no longer propagate: a resource a source stops listing stays indexed forever. The README documents the consequence honestly, but nothing tracks remediation — and this is a behaviour regression to the existing crawler riding along inside a feature PR. Worth a linked follow-up issue at minimum.
  • +113 lines of new config documentation are born stale (config/local-beckn-one-bpp.yaml): they document publishBindingKeys, MANDI_PUBLISH_URL, and "add one entry to the collectors table" — all three retired or deleted by this same PR.
  • The four new config keys are undocumented. publishPipelines, publishEnabled, publishTickIntervalSeconds, publishCatalogOutputDir appear only in a YAML comment, while catalogcrawler/README.md claims to list "the full set of config keys".

Cleanups

  • catalog-preview/mandi-MH.json — ~375 KB of generated output committed (roughly 43% of the diff's line count). Delete and gitignore.
  • Dockerfile.adapter-with-plugins adds COPY tools/ ./tools/, and .gitignore adds /mandi_publish, /markets.json, /mandi_*.json — all for a tool this PR removes. (The apk upgrade zlib alongside .trivyignore is well justified, keep that.)
  • Dead code: asRecords2 (steps.go:761, zero callers); in fetchGet (upstream.go:244-252) both branches of the tokenPlace == "header" if/else are byte-identical.
  • Mandi vocabulary in the generic frame: Outcome.StateCode, the stateErrors counter. Rename to group / groupErrors now, while there's exactly one caller.
  • compiledSchemas (spec.go:461) is an unsynchronised package-level map. Safe today because sweeps are sequential; a data race the moment anything parallelises.
  • publishDiscoverer.Discover makes 1 + N registry HTTP calls every 5-minute tick before the cheap cron/run-log due-check. Flip the order.
  • verdict() reads only answer.Message.Results[0] — a multi-result on_publish would be silently truncated.
  • auth.kind has no enum in the JSON schema, so a typo fails at runtime rather than at load.
  • YAML comments reference dev_docs/... and docs/superpowers/plans/..., both gitignored — unreachable for anyone reviewing this.
  • catalogcrawler now imports its own parent package pkg/plugin/implementation. Works, but it's a latent import-cycle trap.

Suggested path

Merge-blocking from my side: items 1–4 plus the false YAML header. Everything else can be follow-ups.

The cleanest proof that this design does what it sets out to do is to land a second capability — weather — with zero Go changes, either in this PR or immediately after. Right now that isn't possible, because of blocker #1. That's the real acceptance test for what's being built here, and I'd rather we find out now than after the abstraction has hardened.

Good work overall — the bones are solid and this is clearly the right direction.

@manjudr manjudr linked an issue Sep 25, 2026 that may be closed by this pull request
@manjudr

manjudr commented Sep 25, 2026

Copy link
Copy Markdown
Member

Follow-up: concrete fixes for the blockers

My earlier comment was heavier on diagnosis than on remedy. Here are the specific changes I'd propose, so the discussion can be about the shape of the fix rather than about whether the problem is real.

The organising idea for #1: a capability's pre-requisites should be declared in its file and executed by the shared frame — never known to catalogcrawler and never hardcoded in Go. Today that's true structurally (the code lives in catalogpublisher/pipeline) but not behaviourally (there is exactly one hardcoded pre-req, and it's mandatory).


1. Make auth a strategy, not a constant

Today: run.go:245-250 demands tokenUser + tokenSecret whenever an upstream: exists, run.go:286 calls client.Token(...) unconditionally, and upstream.go:88 rejects every auth.kind except tokenExchange. Omitting upstream: isn't an escape hatch — it disables HTTP steps entirely.

Proposed — new pipeline/auth.go:

// Credential is what a request carries, however it was obtained.
type Credential struct {
    Value     string // "" for kind: none
    CarriedAs string // "header" | "query"
    Name      string // "Authorization" | "token" | "x-api-key"
    Prefix    string // "Bearer " | ""
}

// Authenticator runs once, before any step.
type Authenticator interface {
    Prepare(ctx context.Context, c *Client, rc *runContext) (Credential, error)
    RequiredInputs() []string // replaces the hardcoded list in run.go
}

var authKinds = map[string]func(Auth) (Authenticator, error){
    "none":          newNoAuth,
    "apiKey":        newAPIKey,        // reads ${inputs.…}, no round trip
    "basic":         newBasic,
    "tokenExchange": newTokenExchange, // today's code, moved, unchanged
}

run.go then asks the authenticator what it needs instead of asserting it:

if hasUpstream(spec) {
    if resolved[inputBaseURL] == "" { return fmt.Errorf(...) }
    auth, err := authenticatorFor(spec.Upstream.Auth) // empty kind => tokenExchange, back-compat
    if err != nil { return err }
    for _, required := range auth.RequiredInputs() {
        if resolved[required] == "" { return fmt.Errorf("input %q is empty; …", required) }
    }
    client = NewClient(resolved[inputBaseURL]).WithErrorRules(spec.Upstream.Errors)
    cred, err := auth.Prepare(ctx, client, rc)
    if err != nil { return fmt.Errorf("preparing credentials: %w", err) }
    rc = newRunContext(resolved, cred)
}

And one c.applyCredential(req, cred) replaces the hardcoded placements at upstream.go:258 (Authorization: Bearer) and upstream.go:286 (q.Set("token", …)). That's worth doing on its own: it makes auth.token.name live instead of inert, clearing one of the inert-field findings for free, and it deletes the identical-branch if/else at upstream.go:244-252 as a side effect.

Weather then needs no Go:

upstream:
  baseUrl: ${inputs.weatherBaseUrl}
  auth:
    kind: apiKey
    token:
      value: ${inputs.weatherApiKey}
      carriedAs: header
      name: x-api-key

Please also add an enum for auth.kind in pipeline.v1.schema.json — today a typo is a runtime failure, not a load failure.

If "pre-req" means more than auth — e.g. fetch a station list before fanning out — the same shape extends: a prepare: block of ordinary steps that runs before pipeline.steps, with its output addressable as ${prepare.…}. Auth becomes the first built-in prepare step rather than a special case. Worth deciding now which of the two you want, because it changes how much of the above is scaffolding.


2. Make reexchangeOn actually fire

Today: classifiedError returns a plain error for reauth (classify.go:155), so outcomeFor (steps.go:509-512) can't distinguish it and buckets it as transportError → stateErrors++ → refuseWhen: collection.stateErrors > 0 refuses the whole publish. ReexchangeOn is read nowhere.

Proposed — give reauth a sentinel, exactly as emptyResult already has one:

// classify.go
var ErrReauthRequired = errors.New("upstream rejected the token")

case classifyReauth:
    return fmt.Errorf("%s: %w (status %d)", call, ErrReauthRequired, status)
// steps.go
func (r *stepRunner) outcomeFor(step Step, err error) (StepOutcome, bool) {
    class := "transportError"
    switch {
    case errors.Is(err, ErrNoUpstreamData):  class = "emptyResult"
    case errors.Is(err, ErrReauthRequired):  class = "reauth"
    }
    outcome, ok := step.OnError[class]
    return outcome, ok
}

Then, in the step loop, on ErrReauthRequired: re-run Prepare once, refresh rc, retry the step, and only fall through to OnError if it fails again. Cap it at one re-exchange per step (or N per run) so a permanently-bad credential fails fast instead of looping, and count re-exchanges on their own counter — a token flapping on every call should be visible, not silently absorbed.

The classifier rules for [401, 403] already exist and are tested (classify_test.go:22); this just connects them to behaviour.


3. Key the run log on the pipeline, not the capability

run.go:224 writes RecordPipelineRun(ctx, capability, now) into a table whose primary key column is literally named pipeline. The cataloguepublish-<source> convention anticipates two sources for one capability; they'd share a row and starve each other.

-RecordPipelineRun(ctx, capability, now)
+RecordPipelineRun(ctx, opts.Pipeline.Path, now)

plus the matching read side. One line now; a data migration if it ships as-is.


4. Claim the run before doing it, not after

publishSweep's inFlight mutex is per-process, and the marker is written only after a run lasting minutes — so two replicas both see "not due yet", both run, both publish.

Make the existing upsert conditional, and call it before the work with now:

INSERT INTO crawler_pipeline_run (pipeline, last_run_at)
VALUES ($1, $2)                      -- $2 = now
ON CONFLICT (pipeline) DO UPDATE
   SET last_run_at = EXCLUDED.last_run_at
 WHERE crawler_pipeline_run.last_run_at < $3   -- $3 = start of the due window
RETURNING pipeline;

Zero rows returned = another replica owns this window, skip the tick. That alone removes the double-publish. If you also want recovery from a replica that dies mid-run, add started_at + owner columns and treat a claim older than the pipeline timeout as reclaimable — but the conditional upsert is the part I'd not merge without.


5. Correct the mandi YAML header

Lines 12-22 still say upstream, pipeline, discover and registry are INERT, and cite collect.go — deleted in this PR. Those blocks execute now. Since this file is the template the weather author will copy, the header should say what is actually true, e.g.:

# LIVE   upstream              base URL, auth, error classification
# LIVE   pipeline              the step list this run executes, in order
# LIVE   discover, registry    catalogue identity and publish target
#
# Anything listed here is read at runtime. If you add a key that no Go code
# reads, delete it instead -- a rule stated in the file and enforced nowhere
# is worse than no rule.

And on the inert keys generally: I'd delete upstream.guards.*, publish.concurrency and publish.timeout from both the file and the schema rather than leave them declared. additionalProperties: false then makes re-adding them a deliberate act with an implementation attached. spec.go's own doc comment makes this argument better than I can.


None of this changes the architecture — it's the same design, with the two places that quietly assume "the upstream is Agmarknet" opened up. Happy to pair on any of it, or to split #3/#4 into a separate hardening PR if you'd rather keep this one focused on the capability work.

@manjudr

manjudr commented Sep 25, 2026 •

Copy link
Copy Markdown
Member

Forward-looking: what the second API provider will need

Not blocking this PR — mandi needs none of it. Raising it now because these all touch the same two structures (With and the http step), and they are far cheaper to design while there is one caller than after there are four.

Everything below is additive and opt-in. Each proposal is a new optional block; when it is absent, the code path is the one that exists today. I've written the back-compat guarantee explicitly for each, and there's a checklist at the end for what must stay byte-identical for mandi.


A. Pagination — the gap that will bite first

There is no pagination primitive anywhere in the pipeline package (grepped for cursor/offset/nextPage — nothing). Agmarknet returns a state's rows in one response, so this is invisible today. Most third-party APIs are not like that.

Proposed — an optional paginate: block on http.get/http.post only:

- id: rows
  uses: http.get
  with:
    path: /api/observations
    mapping: rows.request.jsonata
  paginate:
    style: cursor            # cursor | offset | page
    nextFrom: meta.next_cursor   # JSONata over the raw response
    as: pageCursor               # exposed to the request mapping as ${page.pageCursor}
    until: nextAbsent            # nextAbsent | emptyPage
    maxPages: 200                # REQUIRED — no unbounded loop, ever

Design notes that matter for safety:

  • maxPages is mandatory in the schema, not defaulted. A provider that always returns a cursor must not be able to loop forever. Hitting the cap is an error, not a silent truncation — a partial catalogue that looks complete is the worst outcome here, and it's the same argument failWhenEmpty already makes.
  • Pages accumulate into the same flat record list the step returns today, so join/derive/dedupe downstream see no shape change at all.
  • Per-page failures reuse the existing onError map — no new error vocabulary. A page that 404s with the provider's "no rows" rule is emptyResult and ends pagination cleanly.
  • First request is unchanged: ${page.*} is empty on page 1, so the existing request mapping produces the exact same URL it does now.

Back-compat: if step.Paginate == nil → call the current single-fetch path verbatim. Guarded at the top, so an unpaginated step cannot reach any new code.


B. Retry and backoff

Also absent. A provider answering 429 + Retry-After can currently only be classified as an error — and given refuseWhen: collection.stateErrors > 0, one throttle refuses the entire publish. Since #A will multiply request counts, these two arrive together in practice.

Proposed — per-provider, under upstream: (not per-step; it's a property of who you're calling):

upstream:
  retry:
    on: [429, 502, 503, 504]
    maxAttempts: 3
    backoff: exponential      # exponential | fixed
    initial: 2s
    max: 30s
    respectRetryAfter: true

The ordering needs to be stated explicitly, because this now interacts with the reauth fix from my previous comment. Three mechanisms can claim the same failed response, so the precedence must be written down and tested:

  1. Classify using the existing upstream.errors rules.
  2. emptyResult → ErrNoUpstreamData, straight to onError. Never retried — "no rows" is an answer, not a failure.
  3. reauth → re-exchange the credential once, retry the request once. Not the backoff path; a dead token is not a busy server.
  4. retryable (matches retry.on) → back off and retry up to maxAttempts.
  5. Anything else → onError as today.

Retry budget should be per run, not just per request, so a provider degrading across 36 states fails loudly instead of quietly turning a 4-minute job into 40.

Back-compat: no retry: block → maxAttempts: 1 → one attempt, current behaviour exactly.


C. Parallel iteration — a real constraint, correctly handled

Correcting something I said elsewhere: step.Concurrency > 1 is refused with an error, not silently ignored (steps.go:156-159), and the comment explains why — the JSONata evaluator keeps built-ins in package-level state, so parallel iteration would race jsonmapper's lock against expr.go's. Refusing is the right call and I'd keep it.

But the constraint is real and will bind: 36 sequential states is fine, 500 stations at one request each is not, and #A multiplies that again.

The fix is not in the step API — it's giving each worker its own evaluator (a pooled Mapper, one per goroutine) so Concurrency can be honoured rather than refused. Worth scoping as its own piece of work; it's the kind of change that wants its own PR and its own race-detector run.

Back-compat: until that lands, the current refusal stays. Nothing changes.


Suggested sequencing

Ordered so each step lands on a smaller blast radius than the next:

  1. Auth strategy (previous comment, Adapter: Mausamgram provider plugin #1) — blocking for a second provider at all.
  2. Reauth wiring (previous comment, feat: add the OAN registry, JSONata mapper and weather provider plugins #2) — small, and #B's precedence rules depend on it existing.
  3. Pagination (#A) — before With grows more callers.
  4. Retry/backoff (#B) — after feat: add the OAN registry, JSONata mapper and weather provider plugins #2, so the precedence is implementable as specified.
  5. Mapper pooling / concurrency (#C) — independent, whenever volume demands it.

Back-compat checklist — must hold after each of the above

  • mandi-price-agmarket.yaml runs unedited and produces a byte-identical catalogue to the one on this branch.
  • Every new block (paginate:, retry:, auth.kind) is optional, and its absence selects today's code path — guarded at entry, not emulated by new code configured to behave like the old.
  • auth.kind unset ⇒ tokenExchange, so no existing file needs touching.
  • maxAttempts unset ⇒ 1. paginate unset ⇒ exactly one request.
  • No new required keys anywhere in pipeline.v1.schema.json.
  • A golden test pinning the current mandi output, added before any of this, so "did we break it" is answered by CI rather than by reading diffs.

That last one I'd genuinely do first — it's cheap, and it converts every item above from "risky refactor" into "change with a safety net."

@manjudr

manjudr commented Sep 25, 2026

Copy link
Copy Markdown
Member

Naming and layout

Checked the new paths against the rest of the repo, the sibling capability plugins, and the Beckn spec. The plugin layout itself is correct — the problems are all spelling, and one of them is in a string that lives in the registry database rather than in this repo, so the fix needs sequencing.


Correct — please don't change

MandiPrice/ matches the established capability-plugin shape exactly, sibling for sibling:

MandiPrice/MandiPrice.go        WeatherObservation/WeatherObservation.go
MandiPrice/prerequisites.go     WeatherObservation/prerequisites.go
MandiPrice/cmd/plugin.go        WeatherObservation/cmd/plugin.go

and it's registered alongside them in install/build-plugins.sh:36, so the .so name matches the id a deployment names in providerSteps. That's the ONIX-expected format and it's right.

The cataloguepublish-<source>/ convention is also deliberate, not random — it's documented in publishpipelines.go:6-7 and the test fixtures already anticipate WeatherObservation/cataloguepublish-imd/ (schedule_test.go:176). Keep the convention. Just fix its spelling.


1. catalogue → catalog

cataloguepublish-agmarket is the only path in the entire repository containing "catalogue". Everything established is American:

Packages catalogcrawler, catalogpublisher, localcatalogblobstore
Identifiers CatalogID ×186, Catalogs ×98, CatalogBlobStore ×40, CatalogPublisher ×34
Beckn spec catalog ×1602 vs catalogue ×6

The new code adds ~140 Catalogue/Catalogues identifiers plus the schemaRef publish.oan/CataloguePipeline/v1 — and then sits inside packages named catalog*. The mixed spelling is already visible within single files.

Proposed: cataloguepublish-* → catalogpublish-*, Catalogues→Catalogs, CataloguePipeline→CatalogPipeline.

⚠️ schemaRef.uses is a contract identifier. publish.oan/CataloguePipeline/v1 is written into every pipeline file and checked at load. Right now exactly one file in the world uses it — this is the cheapest moment this rename will ever be. If it ships first, the options get worse (accept both strings, or bump to /v2).

2. agmarket → agmarknet

The source is Agmarknet (agmarknet.gov.in). The repo agrees — agmarknet appears 730 times across the three casings. agmarket appears 22 times, all of them in or referring to this one folder name. It reads as a typo that got frozen into a path.

3. mandi-price-agmarket.yaml → agmarknet.yaml

The mandi-price- prefix restates the directory it already lives in (MandiPrice/). The repo's own fixtures use the bare form — imd.yaml, minimal.yaml, other.yaml — so the real file is the one file not following the convention its tests describe.

Net result:

- pkg/plugin/implementation/MandiPrice/cataloguepublish-agmarket/mandi-price-agmarket.yaml
+ pkg/plugin/implementation/MandiPrice/catalogpublish-agmarknet/agmarknet.yaml

and weather lands as WeatherObservation/catalogpublish-imd/imd.yaml — which is what schedule_test.go already assumes.

4. catalog-preview/ at the repo root

A new top-level directory holding one generated file (mandi-MH.json, ~375 KB). My first comment argues it shouldn't be committed at all; if a sample output is wanted, it belongs in MandiPrice/catalogpublish-agmarknet/testdata/ next to what produces it — AgricultureFacility/testdata/ is the existing precedent.


Doing this without breaking anything

The one real hazard: the pipeline path is not only in this repo. network.yaml:64-66 says the sweep uses "every ProviderSchema binding in the registry … using that action's mappings as the pipeline path — today agmarknet-live|MandiPrice". So the path is stored in a Sunbird registry record. A repo-only rename leaves the registry pointing at a path that no longer exists, and the sweep stops publishing — at 00:00, quietly.

The second hazard: the embed glob is load-bearing. publishpipelines.go:33 is //go:embed */cataloguepublish-*/…. A folder that stops matching the glob is silently not embedded — it compiles, it starts, and the pipeline is simply absent.

Suggested order:

  1. Add the guard test first. Assert PublishPipelines() returns a non-empty set and that every <Capability>/catalogpublish-*/ directory on disk yields an embedded file. This converts "silently not embedded" into a red CI run. Do this before any rename.
  2. Pin the output. A golden test on the current mandi catalogue, so the rename is provably behaviour-neutral.
  3. git mv the folder and the file (preserves history — please don't delete-and-add).
  4. Update the glob in publishpipelines.go:33 in the same commit as the git mv, never separately.
  5. Sweep the references — I count these: publishpipelines.go:6,7,31,33,56,64,67, publishpipelines_test.go:13,56, catalogcrawler/publishpipelines.go:15, MandiPrice/.../mappings_test.go:25, pipeline/spec.go:346, pipeline/schema/pipeline.v1.schema.json:5, pipeline/schedule_test.go:176,178, pipeline/fixture_test.go:26,40, config/network.yaml:68, config/local-beckn-one-bpp.yaml:174.
  6. Update the Sunbird registry record's mappings value as part of the deploy, not before and not after. This is the step that can't be tested in CI, so it's worth writing down in the runbook.
  7. Do the schemaRef rename (Adapter: Mausamgram provider plugin #1's ⚠️) in the same PR, while one file uses it.

Steps 1–2 are worth doing regardless of whether you take the renames — they're the safety net that makes every other item in my three comments cheaper to act on.

My recommendation: do all of this in this PR, before merge. Every one of these names is load-bearing somewhere — a glob, a registry record, a contract string — and each gets more expensive the moment a second capability copies it. The cost today is one git mv and a sed; after weather, it's a coordinated migration.

@manjudr

manjudr commented Sep 25, 2026

Copy link
Copy Markdown
Member

Structure, duplication, verbosity, logging, error handling

Measured rather than eyeballed. One of these (the tombstone envelope) I'd move up to should-fix-before-merge; the rest are cleanups.


1. Package structure

pkg/plugin/implementation was a namespace directory, and this PR turns it into a package. On development it contains zero .go files — it exists only to hold subdirectories. This PR adds publishpipelines.go there so it can host a //go:embed, and then catalogcrawler/catalogcrawler.go:29 imports its own parent:

import "github.com/beckn-one/beckn-onix/pkg/plugin/implementation"

Nothing is broken today, but every future file added to pkg/plugin/implementation/ is now in a package that one of its own children imports — a cycle waiting for the second commit that touches it.

Suggested: make it a sibling rather than a parent.

- pkg/plugin/implementation/publishpipelines.go          (package implementation)
+ pkg/plugin/implementation/publishpipelines/embedded.go (package publishpipelines)

The embed pattern becomes ../*/catalogpublish-*/…, or the registry moves under catalogpublisher/ where the rest of the publish machinery already lives. Either way the parent directory goes back to being a namespace.

Also: two different files are both named publishpipelines.go — one in that parent package, one in catalogcrawler/ — and one imports the other. They do unrelated things. Suggest naming them for their jobs: embedded.go (the embed registry) and publishsweep.go (the tick loop).

Minor: catalogpublisher/pipeline is 8,541 lines across 10 files. expr.go (611 lines) is a self-contained expression evaluator with no pipeline concepts in it — a natural subpackage if the package keeps growing. Not worth doing now; worth not forgetting.


2. Duplicate logic — tombstone() is a second Beckn envelope ⚠️

sink/request.go:BuildPushBody is the canonical catalog/publish body builder. pipeline/publish.go:265 hand-rolls a second one for retirement, and the two have diverged:

Field BuildPushBody tombstone
context.bppId ✅ from meta.ParticipantID ❌ absent
context.bppUri ✅ from meta.BppURI ❌ absent
schemaTypes / schemaContext ✅ ❌ absent
visibleTo ✅ ❌ absent
catalogType from meta hardcoded REGULAR
updateMode from meta hardcoded MERGE
marshal error returned discarded (body, _ :=)

The identity fields aren't missing because they're unavailable — DiscoverySink holds ParticipantID and BppURI right there (sink.go:31-32). They're missing because the tombstone builder sits in a different package and can't reach them. And config/audit-fields.yaml audits context.bppId on publish actions, so retirements will land in the audit trail with no participant identity on them.

A tombstone is how stale resources are removed — and with the crawl path now on MERGE, retirement is the removal mechanism. It's the wrong thing to have on a divergent code path.

Suggested: delete tombstone() and have the retirement path call BuildPushBody with isActive: false and empty resources, so there is exactly one place that knows the shape of this envelope.

Also still open from my first comment: asRecords2 (steps.go:761) duplicates asRecords and has zero callers; both branches of the tokenPlace if/else in fetchGet (upstream.go:244-252) are byte-identical.


3. Verbosity

Comment density in the new package, against the repo's own baseline:

New Established
mapper.go 35% sink 19%
run.go 31% router 17%
classify.go 31% signer 11%
spec.go 28% publisher 9%
steps.go 22%

Roughly double the repo norm. I want to be careful here: a lot of these comments are genuinely good — they record why a decision was made, which is the kind of comment that earns its keep, and I'd keep most of them.

But the cost is already visible in this PR: the mandi-price-agmarket.yaml header is a ten-line prose block that is now factually false (it declares live blocks INERT and cites a deleted file). Long prose rots, and nothing fails when it does.

Suggested rule of thumb: keep the "why", drop the restatement of what the code plainly says, and where a comment asserts something testable, write the test instead — a test that fails is worth more than a paragraph that lies. spec.go's own doc comment makes exactly this argument about YAML keys; it applies to comments too.


4. Logging — the I/O layer is silent

Log statements per file:

run.go             12        upstream.go        0   ← every outbound HTTP call
steps.go            5        publish.go         0   ← every catalogue publish
publishpipelines.go 5        sink/client.go     0   ← the actual POST
spec.go             1        sink/request.go    0
sink.go             1        sink/batch.go      0
                             catalog.go         0
                             schedule.go        0

Orchestration is well covered — run.go and the crawler log clearly, and run.go:227's ErrorContext on a failed run-log write is exactly the right call. But nothing in the transport layer logs at all, and neither client can: Client (upstream.go:43) has no logger field, and sink.Client is literally struct{ hc *http.Client }.

Concretely, when the 00:00 run fails you get:

step "rows": transport error

with no URL, no status code, no latency, no response size, and no indication which of 36 iterations it was. Diagnosis starts by adding the logging that should have been there.

tombstone publishes (publish.go:180) are entirely silent — nothing records that a catalogue was retired, or that the retirement failed.

Suggested:

  • Thread *slog.Logger into both clients; one Info per outbound call with method, path, status, duration, response bytes, and the forEach item when there is one.
  • One Info per publish and per retirement, with catalog id and outcome.
  • Never log the query string. For this upstream it carries the token — classify.go:141-146 already strips it from error text for exactly this reason, so the rule exists; it just needs to hold for log fields too. Worth a test, since it's the kind of thing a later "let's log the full URL" commit quietly undoes.

5. Error handling

Mostly solid — errors are wrapped with %w and carry the step id. Four specific spots:

  1. steps.go:86 — produced, _ := asRecords(output). The result is used only for the log line "records", len(produced). So a step whose output isn't record-shaped logs records=0 while having actually produced data. That's the log you'd most want to trust during an incident, and it's the one that can silently lie. Either handle the error or log the raw type.
  2. publish.go:266 — body, _ := json.Marshal(...) in tombstone. A marshal failure posts an empty body. Moot if you take the BuildPushBody suggestion above.
  3. mapper.go:41 — go func() { _ = server.Serve(listener) }(). If the loopback mapping server stops serving, every later transform fails with a confusing fetch error instead of the actual cause. At minimum log it; better, capture it so the first mapping failure can report the real reason.
  4. publishpipelines.go:67 — matches, _ := fs.Glob(...). The pattern is a constant, so this genuinely cannot fail — but that makes it an invariant, not an error. A package-level assertion (or the guard test from my naming comment) states it better than a discarded return.

Not findings, for the record: sink/client.go:51's discarded io.ReadAll error is correctly covered by the downstream readable check, and the _ = closeMapper() / _ = os.RemoveAll() deferred cleanups are conventional.


Priority

@manjudr

manjudr commented Sep 25, 2026

Copy link
Copy Markdown
Member

Deep review: mandi-price-agmarket.yaml and mappings/catalog.yaml

Read both line by line. The JSONata in the mapping is skilful — the $number($count(...)) and = '' workarounds are real, measured findings about the pinned jsonata build, and I'd keep every one of those comments. The problems are elsewhere: the two files overlap, and the pipeline file carries a lot of declared-but-dead surface.


A. mandi-price-agmarket.yaml

A1. Configuration that shouldn't ship as-is

A bare public IP is the default upstream (line 90):

baseUrl:
  default: http://34.0.4.235:8080

Plaintext HTTP, a raw IP with no DNS indirection, committed as the default — so a deployment that forgets MANDI_API_URI silently talks to it rather than failing.

What makes this stand out is that the same file argues the opposite case twenty lines earlier, for registryUrl (71-83): "Deliberately NO default … a baked-in hostname that resolves nowhere fails as a DNS error far from here rather than as 'you did not configure the registry'." That reasoning applies verbatim here, and harder — this one does resolve. Drop the default, or at minimum move it behind a hostname and TLS.

Two more defaults that hide misconfiguration:

  • participantId: agmarknet (115) — a default identity. A deployment that forgets it publishes under someone else's name.
  • networkId: oan-dev (119) — a dev network id as the fallback, in a file whose whole purpose is scheduled production publishing. This lands in visibleTo (catalog.yaml:171), so the failure mode is a catalogue silently visible to the wrong network.

A2. Dead declarations

Input Refs in this file Read by Go
registryUrl (84) 0 0
catalogOut (120) 0 0

Both are fully declared — flag, env, comment — and used by nothing. registryUrl carries the file's single longest comment (13 lines of deployment guidance) for an input that is never read.

states (91) is worse than dead: it's flag-only, and its only consumer is the else branch at 216-223, which the file itself documents as "NOT finished … The runner refuses with a clear message." So it's a knob whose only possible effect is to break the run. Either finish it (a split primitive) or delete it — a documented-broken input is an invitation.

retireOld (142) is a one-time migration switch — the comment says so. It's permanently embedded in a capability's pipeline spec, and it's what the weather author will copy. It belongs in an operator runbook, not the spec.

A3. The schedule and the date window may not agree

cron: "0 0 * * *"        # midnight
timezone: Asia/Kolkata
fromDate/toDate: today

At 00:00 IST, "today" is forty seconds old. Agmarknet arrivals are reported during trading hours, so a midnight run for today's date looks like it would fetch a day that has no data yet — every night.

And nothing catches that: emptyResult is recorded to emptyStates and continues (249), refuseWhen only fires on stateErrors (363), and failWhenEmpty guards only the state list (225). So an all-empty run publishes nothing, logs no failure, and records a successful run. The next night it does it again.

I may be wrong about the upstream's semantics — it may serve the previous session under today's date. But the file doesn't say, and that's the point: this is the single most important behavioural assumption in it. Either run at a sensible hour, or default the window to yesterday, or write down why midnight/today is correct.

A4. The failure policy is asymmetric

Lines 179-181 record the lesson beautifully: "Recording 'no rows' as a failure once turned 27 of 36 states into 27 outages." Then line 363 does the mirror image:

refuseWhen: collection.stateErrors > 0

One transport blip in 36 sequential calls discards the other 35 states' data. And because reexchangeOn is inert (my earlier comment), a token expiring mid-run lands in stateErrors too — so the most likely single cause of a failed night is also the one that refuses everything. A threshold (> 3, or > 10%) would match the judgement the file already shows about empty results.

A5. Inert blocks

Still declared, still read by nothing: guards (189-192), auth.token.reexchangeOn (177), auth.token.name (174), publish.concurrency + timeout (355-356). Covered in my earlier comments — listing them here because this is the file they live in.

A6. Dangling references

  • line 24 → dev_docs/plan-crawler-mandi/… — gitignored
  • line 34 → catalogpublisher/pipeline/{schema,reference}/ — reference/ does not exist
  • line 55 → plan-2-low-level-design.md — not in the repo
  • lines 68, 196 → tools/publish/mandi_publish — deleted by this PR
  • lines 17, 158, 199 → collect.go — deleted by this PR

Every pointer a new author would follow is broken.

A7. Spelling, one last time

Line 47 says provider: agmarknet — correct. Line 40 says "it is publish.oan, not publish.agmarket" — the typo. Same file, both spellings, five lines apart.


B. mappings/catalog.yaml

B1. ⚠️ The published envelope has no bppId / bppUri

"context": {
  "action": "catalog/publish",
  "version": "2.0.0",
  "transactionId": …, "messageId": …, "timestamp": …
}

sink.Publish posts a pipeline-built body verbatim (sink.go:88), so this JSONata output is the wire payload. Meanwhile the crawl path's BuildPushBody sets bppId and bppUri from DiscoverySink.ParticipantID/BppURI.

So the two paths this PR exists to unify publish structurally different envelopes, and config/audit-fields.yaml audits context.bppId on publish actions. Combined with the tombstone() builder from my earlier comment, there are now three hand-written versions of this envelope, and two of them omit the participant identity.

B2. The mapping re-implements the pipeline file

Three rules exist in both files, in two different languages:

Rule Pipeline YAML catalog.yaml
drop markets with no commodities exclude 316 $isPublishable 26
honour withoutGeometry: skip exclude 318 $isPublishable 27
order by marketId order 321 $sortedMarkets 31

Both run. They agree today by hand, and nothing enforces that. The YAML exclude also reports a named reason per excluded market — a feature its comment (A-file 268-280) argues for at length — while $isPublishable drops them silently. Whichever one survives, the other should go; I'd keep the declarative exclude/order and delete the JSONata, so the file that states the rule is the file that applies it.

B3. Market sort order is probably wrong

$sort($validMarkets, function($a, $b) { $a.marketId > $b.marketId })    // line 32 — string compare
$sort($mappedCommodities, function($a, $b) { $number($a.code) > $number($b.code) })  // line 83 — numeric

The same file wraps $number(...) for commodity codes and not for market ids. If marketId is numeric-as-string — and $string($row.marketId) at line 48 says it isn't already a string — then "10" < "9" lexically and markets come out misordered. Same class of bug as the $number($count(...)) one the file already documents catching.

B4. district carries an id, not a district

"district": $string($row.districtId),   // line 50

inside $marketObj, whose other fields are marketName and state. districtName exists and is used four lines later for areaName. Anything reading market.district gets "412".

B5. A known-possibly-wrong areaCode ships

Lines 40-45 are admirably honest — "Not verified against the full state list … this produces a wrong areaCode and needs a real lookup table" — and then it ships. An unverified ISO-3166-2 areaCode is worse than an absent one, because consumers will trust the codeScheme. Suggest emitting the ISO area only for a verified allow-list and omitting it otherwise.

B6. Smaller things

  • $toRFC3339Start / $toRFC3339End (16-21) parse dd-MM-yyyy by $substring with no validation — a malformed window yields a syntactically valid, semantically wrong timestamp.
  • updateMode: "MERGE" hardcoded (169): resources dropped upstream are never removed. Same root issue as the crawl path's FULL→MERGE switch.
  • Hardcoded provider claims — historicalDataAvailable: true, historyPeriod: "P1Y", updateFrequency: "P1D", languages: ["en"] (113-118) — asserted about the upstream with nothing verifying them.
  • Three id schemes coexist: catalog:mandi-price:<slug> (142), AGMARKNET-<slug> (145), AGMARKNET-01 (153), and cat-agmarknet-mandi-prices in the pipeline file (367).
  • "catalog:mandi-price:" & _local.catalogSlug is written twice (142, 167) and a third time as identity.catalogId in the pipeline file (333) — which Go really does use (catalog.go:395). Three places that must agree by hand.

What I'd do

Before merge: bppId/bppUri (B1), the IP default and networkId: oan-dev (A1), and either fix or document the midnight/today window (A3).

Before weather copies it: delete the dead inputs and retireOld (A2), collapse the duplicated exclude/order (B2), fix the dangling refs and the header (A6).

Worth a test each: the market sort order (B3) and the $isPublishable-vs-exclude agreement (B2) — both are silent-wrong-answer bugs, which is the kind this file has already been bitten by twice.

@manjudr

manjudr commented Sep 25, 2026

Copy link
Copy Markdown
Member

Reuse audit — what was rebuilt instead of reused

Went through the new code asking one question per component: does this already exist in the repo? Three real cases, one of them a live data race. I've also listed what I checked and cleared, so the list is a conclusion rather than a suspicion.


1. ⚠️⚠️ A second JSONata evaluator — and the mitigation doesn't cover the real exposure

expr.go:56 opens its own engine and guards it with its own package-level mutex:

var evaluating sync.Mutex        // expr.go:47
instance, err := jsonata.OpenLatest()

jsonmapper already has exactly this — jsonmapper.go:166 declares its own var evaluating sync.Mutex around the same library.

The library keeps its built-ins in package-level state, so two mutexes over one global protect nothing from each other. expr.go:41-45 says so outright:

"This lock cannot protect against jsonmapper evaluating at the same time -- that is a different package with a different lock over the same library state. Which is exactly why a step's concurrency is refused above 1."

But refusing concurrency > 1 only prevents parallelism inside one pipeline run. The actual exposure is across subsystems:

  • the publish sweep is "a background job started once at process startup" (network.yaml:60-61), running in the adapter process;
  • that same process serves live traffic through reqmapper → jsonmapper.Transform;
  • schemaversionmediator evaluates JSONata on the request path too.

So a request being mapped at 00:00:03 and the sweep evaluating $count(commodities) = 0 are different goroutines, holding different mutexes, writing the same package frame. Sequential steps inside the pipeline don't help — the concurrency is with the rest of the server.

The irony: run.go:265 already stands up a loopback HTTP server specifically so jsonmapper can fetch this pipeline's mappings. The pipeline is already a jsonmapper client. It then builds a second engine next to it.

Suggested: evaluate bare JSONata through jsonmapper as well, or export one shared evaluator/lock that both packages take. Either way there must be exactly one lock over the one global. Worth a -race test that runs a mapping transform and a pipeline expression concurrently — today I'd expect that to fail.

2. A third catalog/publish envelope builder

Builder bppId/bppUri
sink/request.go BuildPushBody ✅
pipeline/publish.go tombstone() ❌
mappings/catalog.yaml JSONata ❌

Three hand-written versions of one protocol envelope; two omit participant identity, and config/audit-fields.yaml audits context.bppId. Detailed in my previous two comments — repeating here because it belongs on this list.

3. Exclusion and ordering implemented twice

catalog.exclude / catalog.order in the pipeline YAML, and $isPublishable / $sortedMarkets in catalog.yaml. Both run; nothing enforces agreement. Detailed in the previous comment.

4. Dead duplicate

asRecords2 (steps.go:761) duplicates asRecords, zero callers.


Checked and cleared — not duplication

  • store/pipelinerun.go vs store/cursor.go. Different concepts: a per-catalog cursor with a captured envelope, versus a per-pipeline last-run instant. Migration 0008's own comment draws the line correctly.
  • Ad-hoc &http.Client{}. The new code does this twice; the repo already does it in ten places (opapolicychecker, vcvalidator, manifestloader, jsonmapper, …). Consistent with existing practice, not a regression this PR introduced.
  • The two "schedulers" — see below.

Your three questions

"Why a separate scheduler? The crawler has one."

They're not duplicates — they answer different questions:

catalogcrawler/scheduler.go the timer: AddPeriodic, Start/Stop, fixed-interval loops. Knows nothing about cron.
pipeline/schedule.go the due-ness decision: parseCron, dueNow, decideTick, PipelinePathFor. Knows nothing about timers.

The crawler ticks every 5 minutes; schedule.go answers "given this spec's cron and the last-run row, should this pipeline run now?" That split is right, and putting cron semantics next to the spec that declares them is the correct home.

The fair question underneath yours is different: should ~200 lines of cron parser be hand-written at all? There's no cron library in go.mod, so this is all new code — parseCron, parseCronField, cronValue, matches, prevFiring, dayMatches. Having read it, it's good: prevFiring walks backwards (the right question for "was this due and did it run") and skips whole non-matching days, so a yearly expression costs hundreds of iterations, not half a million.

My one concern is DST. It walks backwards minute-by-minute in local time. Asia/Kolkata has no DST so mandi is safe, but the field is per-pipeline and a timezone with DST will hit a local minute that occurs twice (fall back) or not at all (spring forward). I'd add table tests for both transitions before a second pipeline picks a different timezone — or take robfig/cron, which has handled this for a decade.

"Why do we need pipeline.v1.schema.json?"

The honest case for it: it turns a misspelled key into a startup error instead of a rule that silently never fires. additionalProperties: false everywhere means refuseWhn: fails at load rather than at 00:00 against a live upstream. For a config language that operators edit, that's worth real money.

The honest case against it as currently written: it validates keys nothing implements. guards.maxResponseBytes, reexchangeOn, publish.concurrency, publish.timeout all pass validation and are then ignored. A schema that blesses an inert key is worse than no schema — it converts "I'm not sure this works" into "the contract accepted it, so it works."

So: keep the file, and make its scope truthful. Every key it admits must be a key some code reads. That's a deletion exercise, not a redesign — and it's the same argument spec.go's own doc comment already makes.

(Note there are now two schema stacks in-process — this one via santhosh-tekuri/jsonschema for config files, and schemav2validator for Beckn payloads. That's defensible; they validate different things. Flagging it only so it's a decision rather than an accident.)

"Extended schema validation should be enabled if schema validation is enabled"

Agreed, and the new config gets this wrong. config/network.yaml:163-173 — added by this PR:

schemaValidator:
  id: schemav2validator          # ← validation IS on
  config:
    extendedSchema_enabled: "false"   # ← but extended is OFF
    extendedSchema_allowedDomains: "…,openagrinet.github.io"

Every other config in the repo sets it "true" (local-beckn-one-bap.yaml:84, local-beckn-one-bpp.yaml:221). This new file is the outlier.

And it matters specifically for this PR. mappings/catalog.yaml publishes:

"@context": "https://openagrinet.github.io/network-specs/schema/MandiPrice/v0.1/context.jsonld",
"@type": "openagrinet:MandiPrice",

with schemaTypes naming that same context. Extended schema validation is the thing that checks resourceAttributes against the MandiPrice schema pack. With it off, beckn.yaml validates the envelope and nothing validates the payload — so a malformed supportedCommodities, a wrong informationMode, a missing market all publish clean.

The giveaway that this is an oversight rather than a decision: extendedSchema_allowedDomains on line 173 already lists openagrinet.github.io — the exact domain hosting the schema this PR publishes against. The config was prepared for it and then left switched off.

Suggested: set extendedSchema_enabled: "true" in network.yaml, and add a startup check that refuses a config with schemaValidator configured and extendedSchema_enabled: false — so the pairing is enforced rather than remembered. That's your rule, expressed as code.


Where I'd put the effort

The JSONata engine (#1) is the one I'd fix first and independently of everything else in this PR — it's a race in a long-running server process, the stated mitigation doesn't cover it, and the fix is deletion plus reuse of a package this code already depends on. Then extendedSchema_enabled, which is a one-character change that restores validation of the payload this whole PR exists to produce.

Comment thread pkg/plugin/definition/registry.go Outdated
// deployment having to list them: each key is then resolved through
// ProviderRecordLookup, so every gate a single lookup applies still applies.
// Optional, like the other lookups: callers type-assert for it.
type ProviderBindingLister interface {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This turns the registry from a keyed lookup into a list to be walked. Every other use asks ProviderRecord(ctx, bindingKey) with the key coming from the payload.

publishpipelines.go:13 says it itself:

// NOTHING HERE LISTS CAPABILITIES. The registry's publish action names a...

It also scales with the network rather than with this deployment — ProviderRecord is cached per key, a full listing isn't — and it costs three files (this interface, ProviderBindingKeys() in sunbirdRegistry, the sweep that walks it) for discovery the crawler already does from config via NewConfigSource.

If the sweep needs to know what to publish, that reads like the existing crawlmanager.Source seam in internal/source/static.go rather than a new shape on the registry.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Kept the list, deliberately: the registry is meant to be the only place that decides which capabilities publish — a config list would put that decision back in two places. I did narrow it: ProviderBindingKeysServing(ctx, "publish") returns only bindings that serve a publish action (sunbirdRegistry filters its existing listing), so a sweep no longer looks up every binding. Each binding is still judged by the keyed ProviderRecord.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The policy is fine — registry decides what publishes. The objection is the new interface in definition/.

RegistryMetadataLookup.QueryByNetwork(ctx, networkID) already enumerates, and is scoped by network, which this isn't. And the crawler already has crawlmanager.Source.Discover for "what should I work on" — a registry-backed Source lives in the plugin and touches no shared contract.

On the narrowing: ProviderBindingKeysServing still calls searchRecords(..., map[string]eqFilter{}) and filters in Go, so every binding still crosses the wire each tick. Fewer lookups, same fetch.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The concern is definition/registry.go — it's the shared contract every plugin implements, and one feature's discovery need shouldn't grow it. It's now grown twice: ProviderBindingKeysServing was added alongside ProviderBindingKeys rather than replacing it, and both still run the same searchRecords(..., map[string]eqFilter{}).

The policy is fine. Declare the interface where it's consumed instead — Go satisfies them implicitly, so the crawler can write:

type bindingLister interface {
    ProviderBindingKeys(ctx context.Context) ([]string, error)
}

sunbirdRegistry keeps the method, definition/ stays untouched, and no other plugin inherits a concept it doesn't use. publishRegistry in catalogcrawler.go is already a local composite — inlining the signature is the whole change.

// whatever an edited file on the host says, which is not what was reviewed and
// built. The cost is a rebuild to change a pipeline, which is the point.
package implementation

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is the first loose .go file at this level — on development, every entry under pkg/plugin/implementation/ is a plugin directory:

$ git ls-tree --name-only origin/development -- pkg/plugin/implementation/ | grep '\.go$'
(nothing)

It also invents package implementation, which didn't exist before.

What it does is go:embed every capability's pipeline YAML from across the plugin directories. That couples the plugins together at build time, which is the one thing the .so architecture is arranged to avoid — each plugin is independently buildable today, and this makes MandiPrice's files a compile-time dependency of a shared package.

The doc comment also says "Nothing in Go lists the capabilities", but this file embeds them at build time and ProviderBindingKeys() enumerates them at runtime. Both are lists; they're just in different places.

If the pipelines need to reach the binary, embedding them inside the owning plugin package — the way pipeline/mapper.go already serves its embedded mappings — keeps each capability self-contained.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

fixed

@@ -0,0 +1,372 @@
# mandi-price-agmarket.yaml -- the pipeline for the MandiPrice / Agmarknet

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is a third config mechanism. The adapter has two: plugin config under plugins.<name>.config, and JSONata mapping files named by the registry's mappings field. This is neither, yet uses that field.

It also splits cleanly into those two:

  • schedule, inputs, upstream, publish — cron, base URLs, participantId, credentials, timeouts, concurrency. Same kind of thing as discoveryPushUrl and catalogIntervalSeconds; belongs in plugins.crawler.config. Embedded, a cron change is a rebuild.
  • pipeline and catalog — JSONata-shaped, and already half-delegated to mappings/catalog.yaml.

Split that way, pipeline/spec.go, the 1,059-line schema and the expression evaluator have nothing left to do.

Is there something here that can't be plugin config plus mappings?

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fair point, and it's a real design choice rather than a fix. I kept the single file for now, because the goal is that a new capability like weather is one folder plus one registry entry, with no Go and no adapter config change. Splitting it would spread each capability across plugin config and mappings. Happy to discuss it separately, and I'd rather decide it with the weather pipeline in front of us than in the abstract.

Comment thread Dockerfile.adapter-with-plugins Outdated
// catalogue this crawler sends goes. It used to be discovery's /push; a
// config still carrying that is refused rather than left to fail at every
// send (publish bodies at /push, pipelines at .../push/publish).
if strings.HasSuffix(strings.TrimRight(discoveryURL, "/"), "/push") {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

No other code in pkg/ or core/ inspects a URL suffix to decide what an endpoint is — a path is the operator's to choose.

http://provider-adapter:9200/publish is the local compose service name and port, baked into a plugin error. The repo's convention is a placeholder: sunbirdRegistry.go:251 uses e.g. http://<host>:<port>/api/v1.

The refusal also makes this a replacement rather than an extension — a config pointing at /push now fails startup. Both could be supported with a mode key, leaving existing deployments working. cfgDiscoveryURL is still named discoveryPushUrl while requiring a /publish value.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The error (and the README) now show http://:/publish. I kept the /push check on purpose: it refuses a stale discovery address at startup, instead of every send failing later for a reason that's far from the cause.

// Send(ctx, entry, content) contract. It first pushed to a Discovery
// service's /push with UpdateMode FULL -- crawlmanager always resolves a
// catalog's complete current content (via catalog.Resolve), which FULL
// matches. /publish rejects FULL as unsupported, so it now publishes MERGE.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Question: does the crawler need FULL, MERGE, or both?

Asking because the crawler resolves a catalogue's complete current state, which is what FULL meant. Under MERGE a resource the source stops listing is never removed — CONFIG.md says so in this PR: "stays indexed until its catalogue is deactivated."

The stated reason also isn't quite right. discovery-service supports FULL (UpdateModeFull, handled in catalog_repository.go, deletes what the payload omits). What rejects it is this repo's own catalogpublisher — "not yet implemented, rejected with 400".

So if FULL is the mode the crawler actually wants, this reads as a dependency (FULL in catalogpublisher) rather than a limitation to absorb. If MERGE is genuinely right for some sources and FULL for others, that's a per-source choice worth making explicit.

@@ -0,0 +1,628 @@
package pipeline

// build.go runs the `catalog:` block: the part of a pipeline file that turns

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Doc comment says build.go; the file is catalog.go.

Also 12 of the 16 functions here are prefixed catalog — catalogGroup, catalogChunk, catalogCost, catalogRender — in a file called catalog.go, in a package called pipeline. Reads as pipeline.catalogChunk.

The repo's style is the opposite: internal/common/serve.go has serve, resolve, buildRequest, extractAction — no file prefix, most hanging off a receiver. The prefixing here looks like manual namespacing to avoid collisions in a ~30-file flat package, which usually means the behaviour wants a type or the package wants splitting.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

done

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Reopening — the doc comment is fixed, but the naming half is unchanged: still 12 of 16 functions prefixed catalog (catalogGroup, catalogChunk, catalogCost, catalogRender) in pipeline/catalog.go, so calls still read pipeline.catalogChunk.

Not blocking on its own — but it's the symptom worth noting, since the prefixing exists to avoid collisions in a flat ~30-file package.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

naming chages are done

kelvinprabhu added a commit that referenced this pull request Sep 28, 2026
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
@kelvinprabhu
kelvinprabhu force-pushed the feat/117-Mandi-Price-Catalog-Crawler branch from 1f9826f to a05f94d Compare September 28, 2026 11:34
Comment thread pkg/plugin/implementation/catalogpublisher/pipeline/auth.go
@kelvinprabhu

Copy link
Copy Markdown
Collaborator Author

@manjudr thanks for the thorough review. Below is each item with the commit that addressed it. go vet is clean, and go test -race passes for catalogpublisher/..., catalogcrawler/..., MandiPrice/... and sunbirdRegistry. The golden test pins the Mandi output, and every regeneration was a reviewed diff.

Blockers

# Item What changed Commit
B1 Auth hardcoded to token exchange New Authenticator interface (RequiredInputs() / Prepare()) with kinds tokenExchange (the default when unset), none and apiKey. Client.applyCredential replaces both hardcoded placements. The schema has an auth.kind enum and accepts token.value and carriedAs: header. A weather-style kind: apiKey + x-api-key header validates and runs. 57b22aa
B2 reexchangeOn never wired ErrReauthRequired / ReauthError. outcomeFor classifies it as reauth. On a listed status the step re-exchanges once and retries once, capped at 3 per run, counted as reauth, then falls through to onError.reauth. 57b22aa
B3 Run log keyed on capability Every read and write is keyed on the pipeline. It is now the pipeline URL, see "URL-loaded pipelines" below. The crawler's retry budget and give-up mark use the same key. For one release, Run falls back to the old capability-keyed row, so the first tick after deploy does not re-publish a firing that was already served. 57b22aa, 03913d8
B4 No distributed lease Your conditional upsert (ClaimPipelineRun) runs before any fetch. Zero rows means another replica owns the window, and this one stands down. ReleasePipelineRun restores the marker if the run fails, so a transient outage is still retried next tick. Taken literally, the up-front last_run_at = now would have counted a failed run as served. No migration was needed. Tested on real Postgres with 10 concurrent claimants and exactly 1 winner. Not done: crash recovery (started_at/owner). 57b22aa
B5 False YAML header Rewritten; the INERT banners are removed. 8b0d846, 4e60b02

Should-fix

  • Inert fields 57b22aa
    • upstream.guards.*, publish.concurrency and publish.timeout are deleted, and the schema refuses them.
    • reexchangeOn is now live (B2).
    • auth.token.name is now live for query placement and for apiKey.
    • publish.accept / treatAsFailure are dropped, because they only restated the sink's fixed rule ee20a35.
  • Config docs. The four publish keys are documented in catalogcrawler/README.md ee20a35, 03913d8.
  • FULL→MERGE. Not reversible here, because /publish rejects FULL. A follow-up issue will track removing dropped resources.

Reuse audit

  • Two JSONata evaluators. I investigated this rather than changing code. expr_race_test.go runs two independent instances concurrently under -race with no race reported, and reports one as soon as the per-expression lock is removed. The library's only mutable global (staticFrame) is written by init() and RegisterGlobalFunction, and nothing calls RegisterGlobalFunction. The real hazard is one cached *Expression shared across goroutines, and the expr.go lock covers exactly that. 57b22aa
  • Tombstone / envelope. tombstone() is deleted. Retirement goes through BuildPushBody in the sink, and the sink stamps context.bppId/bppUri on every pipeline publish, so the mapping no longer has to know which deployment it runs in. 57b22aa The retirement now carries the file's catalogType, updateMode, visibleTo and schemaTypes, and isActive may only be false — commit pending.
  • Exclude/order duplicated in JSONata. $isPublishable / $sortedMarkets are deleted, so the pipeline file's catalog: block is the only rule and the mapping only renders. ee20a35
  • asRecords2. Removed earlier.

Naming and layout

  • catalogue → catalog everywhere, including the contract publish.oan/CatalogPipeline/v1. agmarket → agmarknet. The file is now agmarknet.yaml. 679278e, 57b22aa
  • The guard test and the golden test went in before the move, as you suggested. 57b22aa
  • The embed is gone. Pipelines are now fetched by URL (below), so pkg/plugin/implementation/ has no .go files again, and catalogcrawler no longer imports its parent package. 9d12469, 03913d8
  • The folder is now MandiPrice/catalogpublish/, with no source suffix, because the provider is the YAML's own name. 9d12469

URL-loaded pipelines 03913d8

  • The registry's publish mappings value is now the pipeline file's https URL. Each run fetches it, with a 30s timeout, a 1 MiB cap and same-host redirects only, then validates it and runs it. If the host is down, the last copy that loaded and validated is used, with a WARN.
  • Guards against redirecting credentials: upstream.baseUrl must be ${inputs.baseUrl}, and every call path must start with /.

Structure, logging, error handling

  • Transport logging 57b22aa
    • pipeline.Client and sink.Client log one line per call: method, path without the query string, status, duration, bytes and the forEach item.
    • The run logs one line per publish and one per retire.
    • Tests assert the token never appears in the logs.
  • Error spots 57b22aa
    • A step output that isn't a collection now logs its type.
    • An unexpected Serve error from the mapping server is logged.
    • The json.Marshal discard went away with tombstone().
  • Generic vocabulary. Outcome.StateCode → Group, and stateErrors/emptyStates → groupErrors/emptyGroups. 57b22aa, ee20a35
  • verdict() now judges every result. Discover resolves only bindings that serve publish, instead of 1 + N lookups. 57b22aa
  • The Dockerfile COPY tools/ and the stale .gitignore lines are removed. 57b22aa

agmarknet.yaml and mappings/catalog.yaml

Item What changed Commit
IP default Removed. allowCleartext is an explicit opt-out that logs a WARN on every run. 8b0d846
participantId / networkId defaults No defaults. A new required: true refuses the run and names the variable to set. ee20a35
Dead inputs states, registryUrl, catalogOut Deleted. The states step is unconditional. ee20a35
Midnight cron vs "today" New yesterday keyword, resolved in schedule.timezone from the run clock. The window defaults to it. ee20a35
refuseWhen > 0 The engine now accepts thresholds; Mandi uses collection.groupErrors > 3. 57b22aa, ee20a35
Dangling doc references Removed from the YAML, spec.go and upstream.go. ee20a35
Market sort was a string compare The mapping no longer sorts. order.by: marketId (numeric) fixes where a state is split into chunks, so a market never moves between catalogIds. The new order.renderBy: [marketName, marketId] sorts each catalog for display. e0f7ed2
district held an id Now districtId plus districtName. ee20a35
Unverified areaCode Explicit Agmarknet → ISO 3166-2 table. Unmapped states get no ISO area. ee20a35
Ids written in three places The engine passes _local.catalogId / $row.resourceId from the file's identity block, so the format is written once. ee20a35
Z timestamps New built-in ${schedule.utcOffset}, so validity times read +05:30. ee20a35
Cron DST Tests found two real bugs: a firing in the spring-forward gap skipped the day, and one in the fall-back hour ran twice. Both fixed. Mandi (IST) was never affected. ee20a35
retireOld in the spec Kept on purpose, default false. It is the migration switch for the old monolithic catalog, and it now carries the deployment's identity. —

Not in this PR

  • Forward-looking items: pagination, retry/backoff and mapper pooling. Now that B2 has landed, the precedence between retry, classification and reauth can be designed.
  • Weather pipeline as the end-to-end "zero Go" proof. Auth no longer blocks it; it will be a separate PR.

Comment thread pkg/plugin/implementation/catalogpublisher/pipeline/interpreter.go
Comment thread pkg/plugin/implementation/catalogpublisher/pipeline/schedule.go
@kelvinprabhu
kelvinprabhu force-pushed the feat/117-Mandi-Price-Catalog-Crawler branch from 46b922f to 7e7e261 Compare September 29, 2026 10:43
@kelvinprabhu

Copy link
Copy Markdown
Collaborator Author

What Changed

Crawler (catalogcrawler)

  • publishDiscoverer.Discover lists the registry's bindings (narrowed to those serving publish when the registry can say so), resolves each, and keeps the ones whose publish action names a pipeline URL.
  • publishsweep.go: a tick every publishTickIntervalSeconds (default 5 min) runs each sanctioned pipeline that is due. Overlap protection: one tick at a time; one pipeline's failure never stops the others.
  • Run log keyed on the pipeline URL (crawler_pipeline_run), so two pipelines for one capability never share a row.
  • Multi-replica safe: a firing is claimed with a conditional upsert (ClaimPipelineRun) before any fetch; a losing replica stands down. A failed run releases the claim (on a context decoupled from the dying one), so the firing is retried rather than lost.
  • Retry budget: 3 attempts per firing for transient failures, then the firing is marked served until the next scheduled one.
  • Permanent faults (bad YAML/schema, non-https URL, literal upstream host, unknown override key, unusable schedule, unimplemented input type) are reported once at WARN, never counted against the budget, and never mark the firing served — so fixing the file lets today's firing run on the next tick.
  • Cancellation (shutdown, deadline) is not a failure: it costs no budget and is not logged at ERROR. A slow upstream (http.Client timeout) still counts as a real failure.
  • A registry value that is not a pipeline URL is logged at ERROR once per value, not on every tick.
  • publishPipelines: "true" enables the sweep; publishBindingKeys is rejected at startup.

Registry contract

  • definition/ is unchanged. The listing methods (ProviderBindingKeys, ProviderBindingKeysServing) are declared as local interfaces in catalogcrawler.go, where they are consumed; sunbirdRegistry satisfies them structurally.
  • A registry that cannot be consulted is an error, never an empty list.

Pipeline engine (catalogpublisher/pipeline)

  • URL-loaded pipelines (source.go): https only (plain http on loopback only), 1 MiB cap, same-host redirects only. The last copy that loaded and validated is kept and used, with a WARN, if the host is briefly down. Mapping refs resolve relative to the pipeline's URL.
  • Credential guards: upstream.baseUrl must be ${inputs.baseUrl} (the address credentials go to comes from the deployment, not the hosted file), and every call path must start with /.
  • Contract: kind: CatalogPipeline, $id publish.oan/CatalogPipeline/v1; validated against the JSON Schema on every load.
  • Schedule: five-field cron resolved against a fixed schedule.utcOffset (e.g. +05:30) — no IANA zone, no zoneinfo database, no DST handling needed. New yesterday date keyword for a midnight firing.
  • Auth kinds: tokenExchange (default), none, apiKey. onError: reauth re-exchanges the token and retries (max 3 per run).
  • Plugin-config overrides, no rebuild: publish.<pipeline>.<input> overrides any non-secret input; publish.<pipeline>.schedule replaces the file's cron. A mistyped key is refused on every tick.
  • Inputs: required: true refuses a run whose input resolves empty.
  • Publish rules: refuseWhen accepts thresholds (collection.groupErrors > 3); publish.accept/treatAsFailure removed (the sink's rule is fixed: only ACCEPTED succeeds). retireOld is fully resolved (isActive: false, updateMode, catalogType, visibleTo, schemaTypes) and sent with the deployment identity.
  • Order: order.by decides where an over-budget group splits (stable ids across runs); order.renderBy sorts each catalog for display.
  • Cancellation: checked between steps and per forEach item; a cancelled run never touches the last good catalogs on disk.
  • HTTP publishing lives in the sink; the engine hands built catalogs to RunOptions.Publisher. upstream: is optional; const step for inline records. One log line per upstream call (method, path without query, status, duration, bytes).

Sink

  • BuildPushBody: action: catalog/publish, schemaTypes in publishDirectives, updateMode: MERGE.
  • Stamps context.bppId/bppUri (from participantId/bppUri) on every pipeline body — the pipeline does not know who it runs as.
  • Client.Push requires HTTP 200 + an ACCEPTED on_publish; DiscoverySink.Publish implements pipeline.Publisher. Push outcomes are audit-logged, correlated by transaction/message id.

Shared HTTP safety (internal/common)

  • One redirect guard, util.RefuseOffHostRedirect, on every credential-bearing client (domain plugins and pipeline): refuses off-host redirects and same-host https→http downgrades, keeps Go's 10-hop limit.
  • A request that ends before the provider answers is a 504, not a 502; a provider that did answer keeps its real status. Attempt counts report attempts actually made. Error bodies are redacted before truncation.

Mandi (MandiPrice/catalogpublish/agmarknet.yaml)

  • Per-state catalogs, one resource per market; participantId and networkId required (no defaults); explicit Agmarknet→ISO 3166-2 table (unmapped states get no area); districtId + districtName; ids from the file's identity block; validity timestamps carry the schedule's offset.

Breaking Changes

  • discoveryPushUrl points to the provider adapter's /publish, e.g. http://provider-adapter:9200/publish. URLs ending in /push are rejected.
  • A binding's publish action mappings is the pipeline's https URL, not a repo path.
  • schedule.timezone → schedule.utcOffset (e.g. "+05:30").
  • Crawled catalogs publish with MERGE; large catalogs may be split across requests.
  • publishBindingKeys removed; every binding with an active publish action runs on a node with publishPipelines: "true".
  • MANDI_PUBLISH_URL → CATALOG_PUBLISH_URL. Mandi now requires MANDI_PARTICIPANT_ID and APP_NETWORK_ID.

Configuration

plugins:
  registry:
    id: sunbirdRegistry
    config:
      url: http://sunbird-registry-service:8081/api/v1
      entity: Participant
      providerEntity: ProviderSchema

  crawler:
    id: catalogcrawler
    config:
      dbDsn: "postgres://crawler:crawler@crawler-db:5432/catalogcrawler?sslmode=disable"
      discoveryPushUrl: "http://provider-adapter:9200/publish"
      participantId: "provider.example"        # published as context.bppId
      bppUri: "https://provider.example"       # published as context.bppUri
      publishPipelines: "true"
      publishEnabled: "false"                  # "true" to post; "false" builds only
      publishTickIntervalSeconds: "300"
      publishCatalogOutputDir: "/tmp/catalogs" # optional: keep built catalogs
      publish.mandi-price.schedule: "0 15 * * *"  # optional override

Registry binding action:

{
  "action": "publish",
  "mappings": "https://raw.githubusercontent.com/<org>/<repo>/<ref>/pkg/plugin/implementation/<Capability>/catalogpublish/<source>.yaml",
  "status": "active"
}

Testing

- go build ./..., go vet ./... clean; full suite passes with -race.
- Engine: schema conformance for every pipeline, golden output for Mandi, URL loading/fallback/guards, cron and utcOffset resolution, overrides, reauth, cancellation (between steps, per forEach item, stale-catalog preservation).
- Crawler: discovery filtering, registry failures, overlap, retry budget, permanent faults (reported once, never served), cancellation vs. upstream timeout, claim/release against real Postgres including a concurrent-claim test.
- Sink: ACCEPTED, PARTIAL, REJECTED, malformed/truncated responses, batching, HTTP errors, schemaTypes, identity stamping.
- Redirect guard: off-host, downgrade, hop limit.
- Live Agmarknet run (earlier on this branch): 10 catalogs / 1,341 markets / 36 states in the first sweep; later sweeps skipped the completed firing.
- Live Agmarknet run (earlier on this branch): 10 catalogs / 1,341 markets / 36 states in the first sweep; later sweeps skipped the completed firing.

Adding a Publishing Capability

1. Create pkg/plugin/implementation/<Capability>/catalogpublish/<source>.yaml and its mappings/.
2. Add a conformance test in the capability's package (the file validates against the schema and every key reaches code).
3. Host the folder (e.g. raw GitHub at a pinned ref).
4. Add the publish action with that URL to the binding's ProviderSchema record.

@manjudr

manjudr commented Sep 29, 2026

Copy link
Copy Markdown
Member

Review — 46b922f

All PR-touched packages pass go test -race -count=1 on go1.26.8 (catalogpublisher, .../pipeline, catalogcrawler, .../sink, .../source, .../store, MandiPrice, MandiPrice/catalogpublish, sunbirdRegistry, internal/common, internal/common/util). The four blockers are genuinely fixed and the logging ask came out well. Detail below, plus one thing I'd want changed before merge and one I'd want a decision on.


1. The JSONata lock comment is wrong, and its test only covers a benign input class

pipeline/expressions.go:36-51 says:

It does NOT need to cover jsonmapper. Measured, not assumed -- expressions_race_test.go runs two independent instances concurrently under -race and reports nothing (...) The library's one mutable global, staticFrame, is written by init() and by RegisterGlobalFunction, which nothing here calls.

Three problems:

a. expressions_race_test.go does not exist. The test is expressions_test.go, and its own header comment also cites the non-existent filename. Both references need fixing regardless of the rest.

b. It contradicts jsonmapper/jsonmapper.go:166, which declares its own var evaluating sync.Mutex and states the opposite conclusion about the same library. Two package-level locks over one shared global, with opposite stated rationales, is the kind of thing that gets "simplified" later by someone who believes whichever comment they read first.

c. staticFrame is not the only shared mutable state. The *Function values bound into it are shared by pointer across every independent OpenLatest() instance, and applyFunction writes to one:

v206/jsonata.go:3521   fn.token = lambda.body.(*ASTNode).Procedure.Value.(string)   // write
v206/jsonata.go:3596   if p.token == "millis" { ... } else if p.token == "now" {    // read

The write is commented "Set token for error reporting", but 3596 uses that same field to select behaviour, so this isn't confined to error text.

I reproduced it. Two goroutines, each with its own OpenLatest() instance, its own compiled expression, and its own lock — i.e. exactly the two-subsystems-with-separate-locks arrangement the comment blesses:

pair result
a + b vs $sum(values) clean — this is the class TestTwoIndependentEvaluatorsDoNotShareState covers
($f := function($x) { $count($x) }; $f(c)) vs $count(c) WARNING: DATA RACE, write 3521 vs read 3596

The trigger is narrow and worth stating precisely, because it bounds the severity:

  • It needs a lambda whose body is a direct call to a builtin in tail position (Procedure.Type == "variable"). A ternary-bodied lambda does not write the token — I checked.
  • It needs both sides to touch the same builtin. $string vs $uppercase does not race.

Against today's repo content this is latent, not live. agmarknet.yaml contains zero lambdas and uses only $count/$exists/$not/$number; the real lambdas in config/mappings/** have ternary or compound bodies. I paired the actual pipeline expression against the actual weather-observation.select.yaml:198 lambda and it ran clean.

So: not a fire. But the comment asserts a general safety property that is false, and the test that allegedly measured it never leaves the safe class. The next person to write function($x) { $count($x) } in a mapping — an entirely ordinary thing to write — turns it on, with nothing in CI to catch it.

Minimum I'd want: fix the two dead filename references, and replace the "measured, not assumed" paragraph with what was actually measured. Better: one lock shared by both packages, since the thing being protected is process-global either way.


2. Review comments — addressed

Comment Status
B1 auth hardcoded to one shape pipeline/auth.go: Authenticator iface, authenticatorFor over tokenExchange/none/apiKey, placement(), WithCredential/applyCredential
B2 re-auth indistinguishable from failure upstream_errors.go: ErrReauthRequired, ReauthError{Call,Status} + Unwrap(), bounded by maxReauthsPerRun
B3 run log keyed on capability run.go: key := opts.Pipeline.URL, documented one-release fallback to the old row
B4 two replicas double-publish store/pipelinerun.go: conditional upsert ON CONFLICT ... WHERE last_run_at < $3 RETURNING; ReleasePipelineRun via context.WithoutCancel
YAML header claimed things the file didn't do Rewritten
catalog-preview/, asRecords2, fetchGet dup branches Deleted
StateCode/stateErrors leaking Mandi into generic code → groupErrors
compiledSchemas unguarded Now under var compiling sync.Mutex
verdict() short-circuited Iterates all results
auth.kind unconstrained Enum in schema
dead guards.* / publish.timeout Deleted
no transport logs Both clients take *slog.Logger; one InfoContext per call
four error-handling spots Fixed
tombstone() hand-rolled Deleted; Retire goes through BuildPushBody
bppId/bppUri missing DiscoverySink.stampIdentity
catalogue→catalog, agmarket→agmarknet, CataloguePipeline→CatalogPipeline Done
parent-package go:embed import Removed; publishsweep.go replaces both old files
LoadLocation/tzdata in a Wolfi container Removed; fixed utcOffset: "+05:30"
new config keys undocumented Documented

Plus golden tests, spec_conformance_test.go, and sink_test.go:420 proving the client logs req.URL.Path and never RawQuery.

3. Review comments — still open

  • Naming convention changed without discussion. I asked to keep cataloguepublish-<source>; this collapses it to plain catalogpublish/, dropping the per-source suffix. schedule_test.go fixtures still anticipate WeatherObservation/…-imd/. Invisible with one source in tree.
  • The acceptance test I named isn't here — landing weather with zero Go changes. Without a second source, "the YAML is the program" is untested against a second instance.
  • @ameersohel45's auth-duplication point is unresolved (pipeline/auth.go vs internal/common/auth.go). "Looking into it", then nothing.
  • config/local-beckn-one-bpp.yaml:173-174 still points at the deleted cataloguepublish-agmarket/mandi-price-agmarket.yaml — the last agmarket strings in the repo.
  • Dangling dev_docs/ refs (gitignored dir): live_test.go:14, local-beckn-one-bpp.yaml:67.

4. Structure

Decomposition is right — pipeline/ executes, catalogcrawler/ schedules, sink/ delivers, sources are data. Three concerns:

a. Remote pipeline fetch has no host allowlist and no signature. pipeline/source.go landed after my last pass, so nobody has reviewed it. A registry record now names an arbitrary https URL that is fetched and interpreted as a program, with its mappings/ resolved relative to it. checkPipelineURL enforces https-or-loopback, fetchPipeline caps 1MB/30s, RefuseOffHostRedirect is set, checkUpstreamIsAnInput constrains baseUrl — all good. But nothing constrains which host, and nothing verifies the bytes. Whoever can write a registry mappings field gets code execution in the adapter's trust domain. We already have the pattern for this (extendedSchema_allowedDomains). This is the one I'd want a decision on before merge.

b. Discovery does the expensive work before the cheap due check. publishDiscoverer.Discover does a registry ProviderRecord lookup and d.resolve(pipelinePath) (the HTTPS fetch) for every binding on every tick; only the resulting targets are then tested against cron/run-log. At defaultPublishTickInterval = 5m that's ~288 registry + HTTPS round trips/day per binding to decide whether a once-daily job should run.

c. The in-repo agmarknet.yaml is a fixture, not a deployable. Tests validate it; production fetches the hosted copy (only non-test references left are catalogcrawler/README.md:74,90). No CI check for drift. The file reads like the source of truth and isn't.

5. Logging

The cleanest part of this PR. Both HTTP clients take *slog.Logger and emit one InfoContext per outbound call with method, path, status, duration, bytes; upstream.go adds the forEach item via withLogItem. run.go carries 18 sites across the run lifecycle. Levels show judgment — Discover logs a registry misconfiguration at Error once and Debug on repeat via the refused map, which is right for something that re-fires every tick.

Verified: both clients log req.URL.Path, never RawQuery; no credential or token value reaches a log call.

One gap: redaction is implemented but only the query-string case has a test. Worth one, given how much of the B1 work moves secrets around.


Summary. Blockers fixed, logging done well. Blocking on §1 — not the lock, the comment: it asserts a safety property I can show is false, backed by a test that never exercises the racing class and a filename that doesn't exist. §4(a) needs a call from us since no one has reviewed source.go.

- Removed default HTTP address for MANDI_API_URI to prevent credential exposure.
- Introduced attempt budget for publishSweep to limit retries on permanent failures.
- Added tests to ensure that permanent failures stop retrying and that budgets reset for new firings.
- Implemented checks to prevent sending credentials over cleartext HTTP connections.
- Enhanced error messages for clarity regarding upstream connection issues.
- Updated schema to include allowCleartext option for upstream configurations.
- Improved handling of case sensitivity in catalogue slugs to prevent overwriting on case-insensitive filesystems.
- Ensured that token exchange and requests do not carry credentials to unintended hosts.
- Introduced `publishsweep_test.go` to test the publish sweep functionality, including configuration validation, registry interactions, and handling of overlapping sweeps.
- Created `embedded.go` to manage embedded publish pipelines, ensuring that capabilities can be added without modifying Go code.
- Implemented `embedded_test.go` to validate that all embedded pipelines conform to their expected contracts and are correctly loaded.
- Renamed and reorganized embedded test files for clarity and consistency.
- Added tests to ensure every embedded pipeline loads correctly and checks for proper imports.
- Introduced a new authentication mechanism in the pipeline, supporting token exchange, API keys, and no authentication.
- Implemented tests for the new authentication strategies, ensuring correct credential handling.
- Added race condition tests for JSONata library to ensure independent evaluators do not share mutable state.
- Created golden tests to pin the exact output of the pipeline, ensuring stability across changes.
- Updated existing tests to reflect changes in the pipeline structure and naming conventions.
…nce input handling, and improve scheduling logic

- Removed the acceptance criteria from the publish spec, as ACCEPTED is the only success state recognized by the publisher.
- Introduced a required flag for inputs, which will cause the run to fail if the input resolves to empty, providing clearer error messages.
- Enhanced the scheduling logic to account for daylight saving time changes, ensuring accurate firing times.
- Updated tests to cover new input requirements and scheduling behavior.
- Cleaned up the codebase by removing unused functions and comments, improving overall readability and maintainability.
- Introduced `spec_conformance_test.go` to validate the YAML to struct contract for the MandiPrice pipeline.
- Implemented tests to ensure all required keys are present and correctly loaded from the pipeline specification.
- Added a new JSON catalog data file for MandiPrice to serve as a reference for testing.
- Removed the deprecated `embedded.go` and `embedded_test.go` files as part of the cleanup process.
- Updated tests to reference pipeline URLs instead of local paths.
- Modified loadRegistryPipeline to validate against the pipeline's URL.
- Introduced RemotePipeline function to handle fetching pipelines from URLs.
- Implemented validation checks for pipeline URLs to ensure they are secure.
- Enhanced error handling for pipeline fetching and validation processes.
- Updated schema documentation to reflect changes in pipeline handling.
Replace schedule.timezone (an IANA zone name resolved via
time.LoadLocation) with schedule.utcOffset (a fixed offset literal,
e.g. "+05:30", parsed into a time.FixedZone). This network's only zone
is IST, which never observes daylight saving, so a fixed offset needs
no zoneinfo database -- removing the time/tzdata dependency and the
now-unreachable DST skip/repeat handling in prevFiring/wallTime.

Cron firing, the yesterday/today date keyword, and ${schedule.utcOffset}
all resolve through the new parseUTCOffset. Schema, spec, agmarknet.yaml,
and all fixtures/tests updated to match.
…est file per source file

catalog.go   -> run.go          (catalog: block building, called only from execute())
publish.go   -> interpreter.go  (publish step)

Test files follow, and every file now pairs 1:1 with its test:
  run_execute_test.go, catalog_test.go       -> run_test.go
  publish_test.go                            -> interpreter_test.go
  schema_test.go                             -> spec_test.go
  concurrency_test.go                        -> spec_test.go (tests Validate)
  inputs_test.go, expr_race_test.go          -> expressions_test.go
  upstream_auth_test.go                      -> upstream_test.go

fixture_test.go stays separate: shared scaffolding used by run/schedule/
source tests, not owned by one file.

No logic changed. gofmt/go vet clean, go test -race passes across
catalogpublisher/..., catalogcrawler/..., MandiPrice/...; golden output
byte-identical.
…context management

- Enhanced `publishSweep` to avoid burning budget on context cancellations and to report permanent faults only once.
- Updated tests to verify behavior on context cancellation and permanent faults.
- Improved `stepRunner` to stop processing items on context cancellation, ensuring accurate error reporting.
- Modified `execute` function to retain last good catalogs during cancellations.
- Introduced utility functions for handling redirects, ensuring credentials are not sent to off-host redirects.
- Added tests for redirect handling and context cancellation scenarios to ensure robustness.
The otel bump off the CVE-affected versions (f1e1795) moved otel/log to
v0.21.0, which drops the package-level KeyValue constructors and replaces
the type with attribute.KeyValue -- EmitAuditLogs was refactored to match.
The sink's publish audit was written against the v0.19 API, so it built on
this branch and failed the moment CI compiled it merged with the base.
- Implemented spec_conformance_test.go to validate the YAML to struct mapping for the MandiPrice pipeline.
- Added tests to ensure all required keys are present and correctly loaded from the pipeline specification.
- Created a golden test data file (catalogs.json) containing sample catalog data for Karnataka and Maharashtra.
- Verified the structure and content of the pipeline specification, including inputs, upstream configuration, and catalog rules.
- Changed package name from `catalogpublish` to `publish` in multiple test files to reflect new organization.
- Updated README and configuration documentation to clarify the removal of `publishTickIntervalSeconds` and its implications.
- Modified `catalogcrawler` to integrate the publish sweep with the catalog-sync ticker, ensuring it runs on every sync tick.
- Refactored scheduler to support after-sync jobs, allowing publish sweeps to run without blocking catalog sync.
- Adjusted tests to align with the new publish pipeline structure and removed references to retired configurations.
- Updated JSON mapping and evaluation logic to ensure thread safety across concurrent evaluations.
- Ensured all references to the old `catalogpublish` path are updated to the new `publish` path in tests and source files.
@kelvinprabhu
kelvinprabhu force-pushed the feat/117-Mandi-Price-Catalog-Crawler branch from 30ef5b1 to 878b074 Compare September 30, 2026 06:10
- Implemented a new concat primitive in the pipeline that allows for appending collections specified in the `with.of` field.
- Updated the interpreter to handle the new concat step and ensure it requires at least two collections.
- Enhanced the schema to validate the concat step and its inputs.
- Added tests to verify the functionality of the concat primitive, ensuring it correctly appends collections and handles edge cases.
- Updated existing mappings and tests to accommodate the new functionality, including handling of relative date inputs like "N days ago".
…ging

- Updated logging levels in catalog crawler to DEBUG for routine operations.
- Enhanced audit logging in sink client to summarize push body without exposing sensitive catalog content.
- Modified credential application to include query keys on all requests, improving security and consistency.
- Introduced telemetry metrics for publish pipeline runs, tracking outcomes and upstream call statistics.
- Implemented structured logging for upstream calls, differentiating between successful and failed requests.
- Added tests to validate telemetry metrics and logging behavior for various scenarios in the catalog publisher pipeline.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Mandi Catalogue Publisher (Network Adaptor)

4 participants