44from datetime import datetime
55from pathlib import Path
66
7+ import polars as pl
8+
79from chameleon_usage .config import SiteConfig , load_config
10+ from chameleon_usage .exceptions import RawTableLoadError , log_raw_table_load_error
11+ from chameleon_usage .extract .dump_db import dump_to_parquet , generate_grant_sql
12+ from chameleon_usage .ingest import load_intervals
813from chameleon_usage .output import compat
14+ from chameleon_usage .pipeline import run_pipeline
15+ from chameleon_usage .schemas import PipelineSpec
916
1017logger = logging .getLogger (__name__ )
1118
1219
13- def _add_shared_args (
14- parser : argparse .ArgumentParser , default : object | None = None
15- ) -> None :
16- parser .add_argument (
17- "--config" ,
18- help = "Path to etc/site.yml (required for process command)" ,
19- default = default ,
20- )
21- parser .add_argument (
22- "--site" ,
23- action = "append" ,
24- help = "Site key from etc/site.yml (repeatable). Defaults to all sites." ,
25- default = default ,
26- )
27- parser .add_argument (
28- "--data-dir" ,
29- help = "Local directory or s3://. Overrides config.data_dir if set." ,
30- default = default ,
20+ def _resolve_sites (args ) -> list [SiteConfig ]:
21+ """Load config, resolve site list and data_dir overrides."""
22+ if not args .config :
23+ raise SystemExit ("Error: --config required" )
24+ sites_config = load_config (args .config )
25+ site_keys = args .site or list (sites_config .keys ())
26+ configs = []
27+ for key in site_keys :
28+ config = sites_config [key ]
29+ if args .data_dir :
30+ config .data_dir = args .data_dir .rstrip ("/" )
31+ configs .append (config )
32+ return configs
33+
34+
35+ # ── Subcommands ──────────────────────────────────────────────────────────
36+
37+
38+ def cmd_extract (args ):
39+ if args .print_grant_sql :
40+ logger .info ("%s" , generate_grant_sql (args .grant_user , args .grant_host ))
41+ return
42+
43+ # Priority: --db-uri > $DATABASE_URI > config.db_uri
44+ db_uri = args .db_uri or os .environ .get ("DATABASE_URI" )
45+
46+ if args .config :
47+ for config in _resolve_sites (args ):
48+ site_db_uri = db_uri or config .db_uri
49+ if not site_db_uri :
50+ raise SystemExit (f"Error: no db_uri for site { config .key } " )
51+ if not config .data_dir :
52+ raise SystemExit (f"Error: no output path for site { config .key } " )
53+ logger .info ("Extracting %s..." , config .key )
54+ dump_to_parquet (site_db_uri , config .data_dir )
55+ elif db_uri :
56+ if not args .data_dir :
57+ raise SystemExit ("Error: --data-dir required when not using --config" )
58+ dump_to_parquet (db_uri , args .data_dir .rstrip ("/" ))
59+ else :
60+ raise SystemExit ("Error: --db-uri, $DATABASE_URI, or --config required" )
61+
62+
63+ def cmd_process (args ):
64+ start = datetime .fromisoformat (args .start_date )
65+ end = datetime .fromisoformat (args .end_date )
66+ spec = PipelineSpec (
67+ group_cols = ("metric" , "resource" , "site" , "collector_type" ),
68+ time_range = (start , end ),
3169 )
3270
71+ output_base = Path (args .output )
72+ output_base .mkdir (parents = True , exist_ok = True )
73+ usage_frames : list [pl .DataFrame ] = []
74+ export_uri = args .export_uri or os .environ .get ("EXPORT_URI" )
3375
34- def parse_args () -> argparse .Namespace :
35- parser = argparse .ArgumentParser ()
36- _add_shared_args (parser )
76+ for config in _resolve_sites (args ):
77+ if not config .data_dir :
78+ raise SystemExit (f"Error: no data_dir for site { config .key } " )
79+
80+ try :
81+ intervals = load_intervals (config .data_dir , spec .time_range ).with_columns (
82+ pl .lit (config .key ).alias ("site" )
83+ )
84+ site_usage = run_pipeline (intervals , spec , resample_interval = args .resample )
85+ except RawTableLoadError as exc :
86+ log_raw_table_load_error (logger , config .key , exc )
87+ continue
88+ except Exception :
89+ logger .exception ("[%s] unhandled exception" , config .key )
90+ raise
3791
92+ output_dir = output_base / config .key
93+ output_dir .mkdir (parents = True , exist_ok = True )
94+ site_usage_df = site_usage .collect ()
95+ site_usage_df .write_parquet (output_dir / "usage.parquet" )
96+ usage_frames .append (site_usage_df )
97+
98+ if export_uri and usage_frames :
99+ combined_usage = pl .concat (usage_frames )
100+ compat_output = compat .to_compat_format (combined_usage )
101+ compat .write_compat_to_db (compat_output , db_uri = export_uri )
102+
103+
104+ # ── Parser ───────────────────────────────────────────────────────────────
105+
106+
107+ def parse_args () -> argparse .Namespace :
108+ parser = argparse .ArgumentParser (
109+ description = "Chameleon cloud resource usage reporting." ,
110+ )
38111 subparsers = parser .add_subparsers (dest = "command" , required = True )
39112
40- extract = subparsers .add_parser ("extract" )
41- _add_shared_args (extract , default = argparse .SUPPRESS )
113+ extract = subparsers .add_parser (
114+ "extract" , help = "Dump database tables to parquet files"
115+ )
116+ extract .add_argument ("--config" , help = "Path to site config YAML." )
117+ extract .add_argument (
118+ "--site" , action = "append" , help = "Site key (repeatable). Defaults to all sites."
119+ )
120+ extract .add_argument (
121+ "--data-dir" , help = "Local directory or s3://. Overrides config data_dir."
122+ )
42123 extract .add_argument (
43124 "--db-uri" ,
44125 help = "Database URI (mysql://user:pass@host:port). Falls back to $DATABASE_URI." ,
45126 )
46-
47- grant_sql = subparsers .add_parser (
48- "print-grant-sql" , help = "Print SQL to grant read access"
127+ extract .add_argument (
128+ "--print-grant-sql" ,
129+ action = "store_true" ,
130+ help = "Print SQL to grant read access and exit." ,
49131 )
50- grant_sql .add_argument ("--user" , default = "usage_exporter" , help = "MySQL username" )
51- grant_sql .add_argument (
52- "--host" , default = "%" , help = "MySQL host patterns (default: %%)"
132+ extract .add_argument (
133+ "--grant-user" , default = "usage_exporter" , help = "MySQL username for grant SQL."
53134 )
135+ extract .add_argument (
136+ "--grant-host" , default = "%%" , help = "MySQL host pattern for grant SQL."
137+ )
138+ extract .set_defaults (func = cmd_extract )
54139
55- process = subparsers .add_parser ("process" )
56- _add_shared_args (process , default = argparse .SUPPRESS )
140+ process = subparsers .add_parser (
141+ "process" , help = "Compute usage timelines from parquet data"
142+ )
143+ process .add_argument ("--config" , required = True , help = "Path to site config YAML." )
144+ process .add_argument (
145+ "--site" , action = "append" , help = "Site key (repeatable). Defaults to all sites."
146+ )
147+ process .add_argument (
148+ "--data-dir" , help = "Local directory or s3://. Overrides config data_dir."
149+ )
57150 process .add_argument (
58151 "--output" ,
59152 required = True ,
@@ -71,130 +164,21 @@ def parse_args() -> argparse.Namespace:
71164 )
72165 process .add_argument (
73166 "--resample" ,
74- help = "Optional resample interval (e.g. 1d, 7d)." ,
167+ help = "Resample interval (e.g. 1h, 1d, 7d)." ,
75168 )
76-
77169 process .add_argument (
78170 "--export-uri" ,
79- help = "Optional DB URI to push output data into. Falls back to env var EXPORT_URI." ,
171+ help = "DB URI to push output data into. Falls back to $ EXPORT_URI." ,
80172 )
173+ process .set_defaults (func = cmd_process )
81174
82175 return parser .parse_args ()
83176
84177
85- def process_site (config : SiteConfig , spec , resample : str ):
86- """Process a site's data through the pipeline. Requires [pipeline] extras."""
87- import polars as pl
88-
89- from chameleon_usage .ingest import load_intervals
90- from chameleon_usage .pipeline import run_pipeline
91-
92- data_dir = config .data_dir
93- if data_dir is None :
94- raise SystemExit (f"Error: no data_dir for site { config .key } " )
95-
96- intervals = load_intervals (data_dir , spec .time_range ).collect ().lazy ()
97-
98- cols = set (intervals .collect_schema ().names ())
99- if "site" not in cols :
100- intervals = intervals .with_columns (pl .lit (config .key ).alias ("site" ))
101-
102- return run_pipeline (intervals , spec , resample_interval = resample )
103-
104-
105178def main () -> None :
106179 logging .basicConfig (level = logging .INFO , format = "%(message)s" )
107180 args = parse_args ()
108-
109- if args .command == "print-grant-sql" :
110- from chameleon_usage .extract .dump_db import generate_grant_sql
111-
112- logger .info ("%s" , generate_grant_sql (args .user , args .host ))
113- return
114-
115- if args .command == "extract" :
116- from chameleon_usage .extract .dump_db import dump_to_parquet
117-
118- # Priority: --db-uri > $DATABASE_URI > config.db_uri
119- db_uri = args .db_uri or os .environ .get ("DATABASE_URI" )
120-
121- if args .config :
122- sites_config = load_config (args .config )
123- site_keys = args .site or list (sites_config .keys ())
124- for site_key in site_keys :
125- config = sites_config [site_key ]
126- site_db_uri = db_uri or config .db_uri
127- if not site_db_uri :
128- raise SystemExit (f"Error: no db_uri for site { site_key } " )
129- # Priority: --data-dir > config.data_dir
130- output_path = args .data_dir or config .data_dir
131- if not output_path :
132- raise SystemExit (f"Error: no output path for site { site_key } " )
133- output_path = output_path .rstrip ("/" )
134- logger .info ("Extracting %s..." , site_key )
135- dump_to_parquet (site_db_uri , output_path )
136- elif db_uri :
137- if not args .data_dir :
138- raise SystemExit ("Error: --data-dir required when not using --config" )
139- dump_to_parquet (db_uri , args .data_dir .rstrip ("/" ))
140- else :
141- raise SystemExit ("Error: --db-uri, $DATABASE_URI, or --config required" )
142- return
143-
144- if args .command == "process" :
145- import polars as pl
146-
147- if not args .config :
148- raise SystemExit ("Error: --config required for process command" )
149-
150- from chameleon_usage .exceptions import (
151- RawTableLoadError ,
152- log_raw_table_load_error ,
153- )
154- from chameleon_usage .schemas import PipelineSpec
155-
156- sites_config = load_config (args .config )
157- site_keys = args .site or list (sites_config .keys ())
158-
159- start = datetime .fromisoformat (args .start_date )
160- end = datetime .fromisoformat (args .end_date )
161- spec = PipelineSpec (
162- group_cols = ("metric" , "resource" , "site" , "collector_type" ),
163- time_range = (start , end ),
164- )
165-
166- output_base = Path (args .output )
167- output_base .mkdir (parents = True , exist_ok = True )
168- usage_frames : list [pl .DataFrame ] = []
169- # Priority: --export-uri > EXPORT_URI
170- export_uri = args .export_uri or os .environ .get ("EXPORT_URI" )
171-
172- for site_key in site_keys :
173- config = sites_config [site_key ]
174- if args .data_dir :
175- config .data_dir = args .data_dir .rstrip ("/" )
176- if not config .data_dir :
177- raise SystemExit (f"Error: no data_dir for site { site_key } " )
178-
179- try :
180- site_usage = process_site (config , spec , args .resample )
181- except RawTableLoadError as exc :
182- log_raw_table_load_error (logger , site_key , exc )
183- continue
184- except Exception :
185- logger .exception ("[%s] unhandled exception" , site_key )
186- raise
187-
188- output_dir = output_base / site_key
189- output_dir .mkdir (parents = True , exist_ok = True )
190- site_usage_df = site_usage .collect ()
191- site_usage_df .write_parquet (output_dir / "usage.parquet" )
192- usage_frames .append (site_usage_df )
193-
194- if export_uri and usage_frames :
195- combined_usage = pl .concat (usage_frames )
196- compat_output = compat .to_compat_format (combined_usage )
197- compat .write_compat_to_db (compat_output , db_uri = export_uri )
181+ args .func (args )
198182
199183
200184if __name__ == "__main__" :
0 commit comments