Skip to content
Merged
Show file tree
Hide file tree
Changes from 3 commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
46 changes: 41 additions & 5 deletions README.md
Original file line number Diff line number Diff line change
@@ -1,8 +1,8 @@
# Dedomena

Consistent research data access for agents and scientists, starting with biology.
OpenAlex, Europe PMC (including PubMed), and EPO patent data share a streaming page
and provenance contract. Queries retain their native source semantics.
Consistent research data access for agents and scientists across biology and finance.
OpenAlex, Europe PMC, EPO, SEC EDGAR, FRED/ALFRED, ECB and World Bank share a
streaming page and provenance contract. Queries retain their native source semantics.

## Install

Expand Down Expand Up @@ -42,18 +42,54 @@ Default caches persist for 24 hours at `~/.cache/dedomena/sources.sqlite3`.
`refresh=True` fetches new data while keeping previous snapshots for offline replay.
`DEDOMENA_SOURCE_STORE` selects a shared store, including across Canary workers.

## Finance sources

~~~python
from dedomena.sources import SEC, FRED, ECB, WorldBank

with SEC() as source: # SEC_USER_AGENT: your organization and contact email
page = source.company_facts(320193) # All concepts, units, and filing vintages
consume(page.records, page.provenance.to_dict())

with FRED() as source: # FRED_API_KEY, free registration
for page in source.release_observations(53): # Up to 500,000 observations/request
consume(page.records, page.provenance.to_dict())
for page in source.observations("GDP", as_of="2020-01-01"):
consume(page.records, page.provenance.to_dict())

with ECB() as source: # No key; currency units per EUR
page = source.fx(["USD", "GBP", "JPY"], frequency="M",
start_period="2025-01", end_period="2025-12")
consume(page.records, page.provenance.to_dict())

with WorldBank() as source: # No key; up to 60 indicators in one query
for page in source.search(["NY.GDP.MKTP.CD", "FP.CPI.TOTL.ZG"],
countries="all", date="1970:2024"):
consume(page.records, page.provenance.to_dict())
~~~

Metadata, units, scaling, missing values, status and filing dates remain native.
FRED supports ALFRED knowledge dates; SEC retains amendments; ECB exposes
revisions; World Bank serves current revised indicators.
See [finance access, throughput and semantics](docs/FINANCE.md).

## Agent CLI

~~~sh
python -m dedomena.sources benchmark openalex 'CRISPR'
python -m dedomena.sources benchmark europepmc 'TITLE:CRISPR'
python -m dedomena.sources search epo 'ta="CRISPR"' --max-pages 1
python -m dedomena.sources quota openalex
python -m dedomena.sources fetch sec 320193
python -m dedomena.sources benchmark fred 53 --operation release
python -m dedomena.sources observations fred GDP --as-of 2020-01-01
python -m dedomena.sources benchmark worldbank 'NY.GDP.MKTP.CD;FP.CPI.TOTL.ZG' --period 1970:2024
python -m dedomena.sources benchmark ecb EXR/M.USD+GBP+JPY.EUR.SP00.A --start-period 2025-01 --end-period 2025-12
~~~

Benchmarks default to one page; `--max-pages` controls acquisition explicitly.
Search streams JSON pages and a final receipt. Source failures emit structured
JSON to stderr with a nonzero exit status. Credentials stay in environment variables.
Collection operations stream JSON pages and a final receipt. Source failures emit structured
JSON to stderr with a nonzero exit status. Credentials stay in environment variables. Offline replay requires no current API key.

See [source limits, usage and Canary integration](docs/SOURCES.md).
Downstream services can expose these same clients to their agents.
Expand Down
6 changes: 5 additions & 1 deletion dedomena/sources/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,8 +4,12 @@
from .openalex import OpenAlex
from .europepmc import EuropePMC
from .epo import EPO
from .sec import SEC
from .fred import FRED
from .ecb import ECB
from .worldbank import WorldBank

__all__ = [
"OpenAlex", "EuropePMC", "EPO", "Page", "Provenance", "Store",
"OpenAlex", "EuropePMC", "EPO", "SEC", "FRED", "ECB", "WorldBank", "Page", "Provenance", "Store",
"SourceError", "BudgetExceeded", "Throttled", "InvalidResponse", "SearchLimitExceeded",
]
130 changes: 95 additions & 35 deletions dedomena/sources/__main__.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,16 +2,19 @@
import argparse
from datetime import date
import json
import os
from pathlib import Path
import sys
import time

from . import EPO, EuropePMC, OpenAlex, SourceError, Store
from . import ECB, EPO, FRED, SEC, EuropePMC, OpenAlex, SourceError, Store, WorldBank


def main(argv=None):
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("action", choices=("search", "fetch", "quota", "usage", "replay", "benchmark"))
parser.add_argument("source", choices=("openalex", "europepmc", "epo"))
operations = ("search", "observations", "release", "catalogue")
parser.add_argument("action", choices=(*operations, "fetch", "quota", "usage", "replay", "benchmark"))
parser.add_argument("source", choices=("openalex", "europepmc", "epo", "sec", "fred", "ecb", "worldbank"))
parser.add_argument("query", nargs="?")
parser.add_argument("--store", help="Shared SQLite snapshot/allowance store")
parser.add_argument("--profile", choices=("discovery", "evidence"), default="discovery")
Expand All @@ -20,78 +23,135 @@ def main(argv=None):
parser.add_argument("--max-pages", type=int)
parser.add_argument("--start-date", type=date.fromisoformat)
parser.add_argument("--end-date", type=date.fromisoformat)
parser.add_argument("--operation", choices=operations, default="search", help="Operation to benchmark")
parser.add_argument("--as-of", help="FRED/ALFRED knowledge date YYYY-MM-DD")
parser.add_argument("--countries", help="World Bank codes joined by semicolons, or all")
parser.add_argument("--period", help="World Bank year/month/quarter or range, e.g. 2020:2025")
parser.add_argument("--start-period", help="ECB ISO date or SDMX period")
parser.add_argument("--end-period", help="ECB ISO date or SDMX period")
parser.add_argument("--updated-after", help="ECB revision delta since ISO timestamp with timezone")
parser.add_argument("--last-n", type=int, help="ECB latest observations per series")
args = parser.parse_args(argv)
operation = args.operation if args.action == "benchmark" else args.action
if args.max_pages is not None and args.max_pages < 1:
parser.error("--max-pages must be positive")
if args.filter and args.source != "openalex":
parser.error("--filter is supported by openalex only")
if args.action in ("fetch", "replay") and not args.query:
if args.as_of and (args.source != "fred" or operation not in ("search", "fetch", "observations")):
parser.error("--as-of requires FRED series search, fetch, or observations")
if (args.countries or args.period) and (args.source != "worldbank" or operation not in ("search", "observations")):
parser.error("--countries and --period require World Bank search or observations")
if (args.start_period or args.end_period or args.updated_after or args.last_n is not None) and (args.source != "ecb" or operation not in ("search", "observations")):
parser.error("SDMX period/delta/last-n options require ECB search or observations")
if args.action != "benchmark" and args.operation != "search":
parser.error("--operation selects the benchmark operation only")
if operation == "release" and args.source != "fred":
parser.error("release requires fred")
if operation == "catalogue" and args.source not in ("ecb", "worldbank"):
parser.error("catalogue requires ecb or worldbank")
if operation == "observations" and args.source not in ("fred", "ecb", "worldbank"):
parser.error("observations requires fred, ecb, or worldbank")
if operation in ("fetch", "replay", "release", "observations") and not args.query:
parser.error("this action needs an identifier")
if args.action in ("search", "benchmark") and not args.query and not (args.source == "openalex" and args.filter):
if operation == "search" and not args.query and not (args.source == "openalex" and args.filter):
parser.error("search needs a query, or an OpenAlex filter")
if bool(args.start_date) != bool(args.end_date) or (args.start_date and args.source != "epo"):
parser.error("date partitions require both dates and the epo source")
store = Store(args.store) if args.store else None
constructors = {"openalex": OpenAlex, "europepmc": EuropePMC, "epo": EPO}
kwargs = {"store": store} if store else {}
if args.source == "openalex":
kwargs["profile"] = args.profile
elif args.source == "europepmc":
kwargs["result_type"] = "core" if args.profile == "evidence" else "lite"
source = constructors[args.source](**kwargs)
store = None

def emit(value):
print(json.dumps(value, ensure_ascii=False, allow_nan=False))

def pages(source):
if operation == "release":
if not args.query.isascii() or not args.query.isdecimal():
raise ValueError("release identifier must be a positive integer")
return source.release_observations(int(args.query), refresh=args.refresh)
if operation == "catalogue":
return iter((source.dataflows(refresh=args.refresh),)) if args.source == "ecb" else source.indicators(refresh=args.refresh)
if args.source == "fred":
method = source.observations if operation == "observations" else source.search
return method(args.query, as_of=args.as_of, refresh=args.refresh)
if args.source == "worldbank":
options = {"refresh": args.refresh}
if args.countries:
options["countries"] = args.countries
if args.period:
options["date"] = args.period
return source.search(args.query, **options)
if args.source == "ecb":
flow, key = source._query(args.query)
return iter((source.observations(flow, key, start_period=args.start_period,
end_period=args.end_period, updated_after=args.updated_after,
last_n=args.last_n, refresh=args.refresh),))
if args.source == "openalex":
return source.search(args.query, filter=args.filter, refresh=args.refresh)
if args.source == "epo" and args.start_date:
return source.search_partitioned(args.query, args.start_date, args.end_date,
refresh=args.refresh,
biblio=args.profile == "evidence")
if args.source == "epo":
return source.search(args.query, refresh=args.refresh, biblio=args.profile == "evidence")
return source.search(args.query, refresh=args.refresh)

try:
with source:
if args.action == "replay":
store = Store(args.store or os.environ.get("DEDOMENA_SOURCE_STORE",
str(Path.home() / ".cache" / "dedomena" / "sources.sqlite3")))
raw = store.replay(args.query)
if raw.provenance.source != args.source:
raise SourceError("snapshot belongs to a different source")
emit({"provenance": raw.provenance.to_dict(), "body": raw.body.decode("utf-8")})
return 0
store = Store(args.store) if args.store else None
constructors = {"openalex": OpenAlex, "europepmc": EuropePMC, "epo": EPO,
"sec": SEC, "fred": FRED, "ecb": ECB, "worldbank": WorldBank}
kwargs = {"store": store} if store else {}
if args.source == "openalex":
kwargs["profile"] = args.profile
elif args.source == "europepmc":
kwargs["result_type"] = "core" if args.profile == "evidence" else "lite"
with constructors[args.source](**kwargs) as source:
if args.action == "quota":
emit(source.quota())
elif args.action == "usage":
emit(source.usage())
elif args.action == "replay":
raw = source.store.replay(args.query)
emit({"provenance": raw.provenance.to_dict(), "body": raw.body.decode("utf-8")})
elif args.action == "fetch":
emit(source.fetch(args.query, refresh=args.refresh).to_dict())
options = {"refresh": args.refresh}
if args.source == "fred":
options["as_of"] = args.as_of
emit(source.fetch(args.query, **options).to_dict())
else:
if args.source == "openalex":
pages = source.search(args.query, filter=args.filter, refresh=args.refresh)
elif args.source == "epo" and args.start_date:
pages = source.search_partitioned(args.query, args.start_date, args.end_date,
refresh=args.refresh,
biblio=args.profile == "evidence")
elif args.source == "epo":
pages = source.search(args.query, refresh=args.refresh, biblio=args.profile == "evidence")
else:
pages = source.search(args.query, refresh=args.refresh)
limit = args.max_pages if args.max_pages else (1 if args.action == "benchmark" else None)
started = time.perf_counter()
stream = pages(source)
limit = args.max_pages or (1 if args.action == "benchmark" else None)
totals = {"pages": 0, "records": 0, "credits_used": 0,
"wire_bytes": 0, "decoded_bytes": 0, "cache_hits": 0}
last = None
for page in pages:
for page in stream:
last = page
totals["pages"] += 1
totals["records"] += len(page.records)
totals["credits_used"] += page.credits_used
totals["wire_bytes"] += page.wire_bytes
totals["decoded_bytes"] += page.decoded_bytes
totals["cache_hits"] += int(page.cache_hit)
if args.action == "search":
if args.action != "benchmark":
emit(page.to_dict())
if limit and totals["pages"] >= limit:
break
emit({"source": args.source, **totals,
emit({"source": args.source, "operation": operation, **totals,
"elapsed_seconds": round(time.perf_counter() - started, 3),
"complete": bool(last and last.complete),
"next_cursor": last.next_cursor if last else None,
"records_seen": last.records_seen if last else 0,
"last_parameters": dict(last.provenance.parameters) if last else {},
"usage": source.usage()})
except (SourceError, ValueError, KeyError) as error:
# KeyError is a snapshot ID, never a raw upstream body or credential.
message = "snapshot not found" if isinstance(error, KeyError) else str(error)
print(json.dumps({"source": args.source, "error": type(error).__name__,
"message": message, "complete": False}), file=sys.stderr)
emit_error = {"source": args.source, "error": type(error).__name__,
"message": message, "complete": False}
print(json.dumps(emit_error), file=sys.stderr)
return 2
finally:
if store:
Expand Down
Loading
Loading