Feat/117 mandi price catalog crawler - #54
kelvinprabhu wants to merge 67 commits into
Conversation
🛡️ Trivy security scan (CRITICAL,HIGH,MEDIUM,LOW)Go dependenciesNo findings at CRITICAL,HIGH,MEDIUM,LOW. Container imageNo findings at CRITICAL,HIGH,MEDIUM,LOW. |
|
📊 Test Coverage: ✅ Passed — 85% of changed lines covered, min 80% |
manjudr
left a comment
There was a problem hiding this comment.
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 | |
| 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.yamlheader (lines 12–22) is now false. It declaresupstream,pipeline,discoverandregistryas INERT and points atcollect.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 hardcodesmaxResponseBytesas a const inupstream.go),auth.token.reexchangeOn, andauth.token.name(the header path hardcodesAuthorization: Bearer, the POST path hardcodesq.Set("token", ...)), pluspublish.concurrencyandpublish.timeout. This is precisely the failure modespec.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 documentpublishBindingKeys,MANDI_PUBLISH_URL, and "add one entry to thecollectorstable" — all three retired or deleted by this same PR. - The four new config keys are undocumented.
publishPipelines,publishEnabled,publishTickIntervalSeconds,publishCatalogOutputDirappear only in a YAML comment, whilecatalogcrawler/README.mdclaims 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-pluginsaddsCOPY tools/ ./tools/, and.gitignoreadds/mandi_publish,/markets.json,/mandi_*.json— all for a tool this PR removes. (Theapk upgrade zlibalongside.trivyignoreis well justified, keep that.)- Dead code:
asRecords2(steps.go:761, zero callers); infetchGet(upstream.go:244-252) both branches of thetokenPlace == "header"if/else are byte-identical. - Mandi vocabulary in the generic frame:
Outcome.StateCode, thestateErrorscounter. Rename togroup/groupErrorsnow, 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.Discovermakes 1 + N registry HTTP calls every 5-minute tick before the cheap cron/run-log due-check. Flip the order.verdict()reads onlyanswer.Message.Results[0]— a multi-resulton_publishwould be silently truncated.auth.kindhas no enum in the JSON schema, so a typo fails at runtime rather than at load.- YAML comments reference
dev_docs/...anddocs/superpowers/plans/..., both gitignored — unreachable for anyone reviewing this. catalogcrawlernow imports its own parent packagepkg/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.
Follow-up: concrete fixes for the blockersMy 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 1. Make auth a strategy, not a constantToday: Proposed — new // 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
}
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 Weather then needs no Go: upstream:
baseUrl: ${inputs.weatherBaseUrl}
auth:
kind: apiKey
token:
value: ${inputs.weatherApiKey}
carriedAs: header
name: x-api-keyPlease also add an If "pre-req" means more than auth — e.g. fetch a station list before fanning out — the same shape extends: a 2. Make
|
Forward-looking: what the second API provider will needNot blocking this PR — mandi needs none of it. Raising it now because these all touch the same two structures ( 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 firstThere 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 - 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, everDesign notes that matter for safety:
Back-compat: B. Retry and backoffAlso absent. A provider answering Proposed — per-provider, under upstream:
retry:
on: [429, 502, 503, 504]
maxAttempts: 3
backoff: exponential # exponential | fixed
initial: 2s
max: 30s
respectRetryAfter: trueThe 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:
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 C. Parallel iteration — a real constraint, correctly handledCorrecting something I said elsewhere: 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 Back-compat: until that lands, the current refusal stays. Nothing changes. Suggested sequencingOrdered so each step lands on a smaller blast radius than the next:
Back-compat checklist — must hold after each of the above
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." |
Naming and layoutChecked 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
and it's registered alongside them in The 1.
|
| 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:
- 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. - Pin the output. A golden test on the current mandi catalogue, so the rename is provably behaviour-neutral.
git mvthe folder and the file (preserves history — please don't delete-and-add).- Update the glob in
publishpipelines.go:33in the same commit as thegit mv, never separately. - 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. - Update the Sunbird registry record's
mappingsvalue 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. - Do the
schemaRefrename (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.
Structure, duplication, verbosity, logging, error handlingMeasured 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
import "github.com/beckn-one/beckn-onix/pkg/plugin/implementation"Nothing is broken today, but every future file added to Suggested: make it a sibling rather than a parent. The embed pattern becomes Also: two different files are both named Minor: 2. Duplicate logic —
|
| 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.Loggerinto both clients; oneInfoper outbound call with method, path, status, duration, response bytes, and theforEachitem when there is one. - One
Infoper publish and per retirement, with catalog id and outcome. - Never log the query string. For this upstream it carries the token —
classify.go:141-146already 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:
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 logsrecords=0while 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.publish.go:266—body, _ := json.Marshal(...)intombstone. A marshal failure posts an empty body. Moot if you take theBuildPushBodysuggestion above.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.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
- Before merge: the
tombstoneenvelope (feat: add the OAN registry, JSONata mapper and weather provider plugins #2) — it's on the removal path, and removal is now load-bearing because of the FULL→MERGE switch. - Before the next capability: transport logging (Docker compose for experience, network and provider layers #4) and the package restructure (Adapter: Mausamgram provider plugin #1) — both get more expensive once a second pipeline exists.
- Whenever: dead code, verbosity, the three remaining error spots.
Deep review:
|
| 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: todayAt 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 > 0One 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) parsedd-MM-yyyyby$substringwith 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), andcat-agmarknet-mandi-pricesin the pipeline file (367). "catalog:mandi-price:" & _local.catalogSlugis written twice (142, 167) and a third time asidentity.catalogIdin 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.
Reuse audit — what was rebuilt instead of reusedWent 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.
|
| 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.govsstore/cursor.go. Different concepts: a per-catalog cursor with a captured envelope, versus a per-pipeline last-run instant. Migration0008'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.
| // 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 { |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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 | ||
|
|
There was a problem hiding this comment.
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.
| @@ -0,0 +1,372 @@ | |||
| # mandi-price-agmarket.yaml -- the pipeline for the MandiPrice / Agmarknet | |||
There was a problem hiding this comment.
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 asdiscoveryPushUrlandcatalogIntervalSeconds; belongs inplugins.crawler.config. Embedded, a cron change is a rebuild.pipelineandcatalog— JSONata-shaped, and already half-delegated tomappings/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?
There was a problem hiding this comment.
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.
| // 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") { |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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. |
There was a problem hiding this comment.
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 | |||
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
naming chages are done
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
1f9826f to
a05f94d
Compare
|
@manjudr thanks for the thorough review. Below is each item with the commit that addressed it. Blockers
Should-fix
Reuse audit
Naming and layout
URL-loaded pipelines 03913d8
Structure, logging, error handling
|
| 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.
46b922f to
7e7e261
Compare
What ChangedCrawler (
|
Review —
|
| 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.
$stringvs$uppercasedoes 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 plaincatalogpublish/, dropping the per-source suffix.schedule_test.gofixtures still anticipateWeatherObservation/…-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.govsinternal/common/auth.go). "Looking into it", then nothing. config/local-beckn-one-bpp.yaml:173-174still points at the deletedcataloguepublish-agmarket/mandi-price-agmarket.yaml— the lastagmarketstrings 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.
…improve token management
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.
30ef5b1 to
878b074
Compare
…verId, removing deprecated bppId and bppUri
- 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.
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
buildPublishSource/publishDiscoverer.Discover.ProviderSchemabindings with an activepublishaction and an embedded pipeline.publishpipelines.go.publishPipelines: "true"enables publishing.publishBindingKeysis now rejected.Sink
BuildPushBodyusesaction: catalog/publish.schemaTypestopublishDirectives.MERGEinstead ofFULL.Client.Pushrequires HTTP200+ACCEPTEDon_publishverdict.DiscoverySink.Publishimplementingpipeline.Publisher.Pipeline
RunOptions.Publisherhandles it.catalogpublisher/catalogpublish; mapper helpers moved topipeline/mapper.go.upstream:is optional; pipelines without it skip credentials/token exchange.conststep for inline records.Registry & Embedding
ProviderBindingLister/ProviderBindingKeys.pkg/plugin/implementation/*/cataloguepublish-*/pipeline_files.goand the crawlercollectorstable.CATALOG_PUBLISH_URL.Breaking Changes
discoveryPushUrlnow points to the provider adapter/publish, e.g.:URLs ending in
/pushare rejected.Crawled catalogues use
MERGE; large catalogues may be split into multiple requests.publishBindingKeysis removed; usepublishPipelines: "true".Every registry binding with a
publishaction runs on the node; there is no binding filter.MANDI_PUBLISH_URL→CATALOG_PUBLISH_URL.Configuration
A registry binding enables a pipeline with:
{ "action": "publish", "mappings": "pkg/plugin/implementation/<Capability>/cataloguepublish-<source>/<pipeline>.yaml", "status": "active" }Testing
go build ./...andgo vetpass.-race.ACCEPTED,PARTIAL,REJECTED, malformed/truncated responses, batching, HTTP errors, andschemaTypes.Adding a Publishing Capability
Create:
pkg/plugin/implementation/<Capability>/cataloguepublish-<source>/Add the pipeline YAML and
mappings/.Run:
go test ./pkg/plugin/implementation/Add the
publishaction to theProviderSchemaregistry record.Rebuild.
The registry is now the single source of truth for publishing capability selection.