Compare commits

...

29 Commits

Author SHA1 Message Date
b49153bd38 some refactoring, fix stringer func 2026-08-08 22:39:37 +02:00
a052c70baa fix cpu percentage element count 2026-08-08 22:39:12 +02:00
35bd685cfe not needed 2026-08-07 23:00:58 +02:00
ce5d7f51c6 streamline table api 2026-08-07 22:59:33 +02:00
59331fc721 satisfy linter 2026-07-17 10:13:32 +02:00
aba5985a2b fix too wide tables when in !tty mode 2026-07-17 10:09:05 +02:00
281499166b fix header count on task ls 2026-07-17 10:08:39 +02:00
744cdcb733 fix word wrapping 2026-07-17 09:52:55 +02:00
6f3ca3243a fix crash 2026-07-17 08:55:48 +02:00
T. von Dein
522898f187 add index copy (aka reindex), fix typoe indicies to indices (#100) 2026-07-16 23:21:45 +02:00
T. von Dein
59cdc03ea2 hide ds backend indices as any other hidden indices (#99) 2026-07-16 21:46:50 +02:00
T. von Dein
d69984b01f satisfy linter (#98) 2026-07-15 23:57:02 +02:00
T. von Dein
00c5d79794 use mapmap.Slicer to filter index templates ls and index ls (#96) 2026-07-15 23:47:19 +02:00
T. von Dein
55a1005857 add index template completion, index pattern filter and name filter (#95) 2026-07-15 14:52:30 +02:00
T. von Dein
2f13e08536 fix #93 crash index template, check nil ptr (#94) 2026-07-15 13:54:39 +02:00
T. von Dein
97a509da9e adopt latest golang enhancements (#91) 2026-07-13 14:33:41 +02:00
c5a21bd677 add esql sample 2026-07-10 15:14:23 +02:00
T. von Dein
f61fda0697 add feature esql search (#90) 2026-07-10 15:07:57 +02:00
T. von Dein
15dbf7ada0 fix api header spec (#89) 2026-07-10 13:58:56 +02:00
T. von Dein
311a8474d6 feature/node-usage (#88) 2026-07-10 13:37:59 +02:00
5644046c78 api repl: print output directly if not json 2026-07-10 11:28:25 +02:00
T. von Dein
9a42b86d62 add support for human readable /_cat endpoints: 'api repl -H' (#87) 2026-07-10 11:21:45 +02:00
T. von Dein
53abe666e8 fix 'node show' and 'shard explain' crashes (#86) 2026-07-09 13:41:07 +02:00
1e9e419e9b fix crash in ByteSize conversion 2026-07-08 15:07:01 +02:00
T. von Dein
c128e32ca1 add more important node stats (#82) 2026-07-08 15:03:31 +02:00
2969522913 rename var 2026-07-08 15:00:21 +02:00
51480e6b35 add 'index du <index>' 2026-07-08 14:58:07 +02:00
99bf703997 add table.AddRowLate() to add rows after sorting 2026-07-08 14:57:46 +02:00
T. von Dein
76ad6c7afa remove workaround to load config in root.Before() (#83)
see: https://github.com/urfave/cli/issues/2348
2026-07-08 14:13:47 +02:00
61 changed files with 2217 additions and 416 deletions

View File

@@ -67,10 +67,10 @@ test: clean buildlocal
testlint: test lint testlint: test lint
lint: lint-basic:
golangci-lint run --enable-only errcheck,govet,ineffassign,staticcheck,unused golangci-lint run --enable-only errcheck,govet,ineffassign,staticcheck,unused
lint-full: lint:
golangci-lint run --show-stats=false golangci-lint run --show-stats=false
testfuzzy: clean testfuzzy: clean
@@ -110,3 +110,12 @@ profile: buildlocal
./esctl api ls --profile-file cpu.profile ./esctl api ls --profile-file cpu.profile
go tool pprof -text esctl cpu.profile go tool pprof -text esctl cpu.profile
go tool pprof --http localhost:8888 ./esctl cpu.profile go tool pprof --http localhost:8888 ./esctl cpu.profile
docker-up:
make -C t up
docker-waitup:
make -C t up
docker-down:
make -C t down

View File

@@ -29,6 +29,7 @@ Features:
logical condition (OR, AND), use PIT, limit datetime (ES date math logical condition (OR, AND), use PIT, limit datetime (ES date math
can be used), etc. It is however not yet possible to create can be used), etc. It is however not yet possible to create
recursive searches like: `(cond1 AND cond2) OR (cond3 OR cond4)`. recursive searches like: `(cond1 AND cond2) OR (cond3 OR cond4)`.
- Search using ES|QL language: `esctl searchql`.
- Cross cluster replication (ccr): view, pause, resume, delete - Cross cluster replication (ccr): view, pause, resume, delete
replication. You can also manage follower configuration. replication. You can also manage follower configuration.
- Index management: manage aliases, create, modify, delete indices, - Index management: manage aliases, create, modify, delete indices,
@@ -215,7 +216,7 @@ Pending Tasks 0
Nodes 3 Nodes 3
Red Indices 0 Red Indices 0
Long Running Tasks 2 Long Running Tasks 2
Indicies 221 indices 221
Docs 9854777 Docs 9854777
Total Size 3.4 GB Total Size 3.4 GB
Total Queries 4210411 Total Queries 4210411
@@ -352,6 +353,30 @@ $ esctl search -i foo* -F title=zeitbuchung message=pause | jq
} }
``` ```
You can also search using [ES|QL](https://www.elastic.co/docs/reference/query-languages/esql/esql-getting-started):
```console
$ esctl searchql "from hyperdrive | sort @timestamp | limit 5"
@TIMESTAMP MESSAGE TAG
2026-06-24T08:24:56.000Z arosu loop
2026-06-24T08:24:58.000Z hami loop
2026-06-24T08:24:59.000Z ishininu loop
2026-06-24T08:25:00.000Z uyomoruron loop
2026-06-24T08:25:02.000Z ishimime loop
```
There are several output modes (json, yaml, csv), to get esql output as CSV:
```console
$ esctl searchql "from hyperdrive | sort @timestamp | limit 5" -o csv
@timestamp,message,tag
2026-06-24T08:24:56.000Z,arosu,loop
2026-06-24T08:24:58.000Z,hami,loop
2026-06-24T08:24:59.000Z,ishininu,loop
2026-06-24T08:25:00.000Z,uyomoruron,loop
2026-06-24T08:25:02.000Z,ishimime,loop
```
To check which field mappings are available for an index: To check which field mappings are available for an index:
```console ```console
$ esctl index show foo2 $ esctl index show foo2
@@ -509,7 +534,7 @@ cluster - manage cluster[s]
allocate-empty-primary - allocate an empty primary shard to a node allocate-empty-primary - allocate an empty primary shard to a node
allocate-stale-primary - allocate a stale primary shard to a node allocate-stale-primary - allocate a stale primary shard to a node
datastream - manage data streams datastream - manage data streams
list - list indicies list - list indices
show - show details about an data stream show - show details about an data stream
create - create a new data stream create - create a new data stream
delete - delete a data stream delete - delete a data stream
@@ -529,8 +554,8 @@ ilm - manage index lifecycle
list - list index rollover config list - list index rollover config
show - show rollover forecast over all indices show - show rollover forecast over all indices
explain - explain ilm condition of an index explain - explain ilm condition of an index
index - manage indicies index - manage indices
list - list indicies list - list indices
show - show details about an index show - show details about an index
create - create a new index create - create a new index
update - update an index update - update an index
@@ -538,6 +563,8 @@ index - manage indicies
close - close an index close - close an index
fields - show info about field capabilities fields - show info about field capabilities
ilm - show ilm status ilm - show ilm status
du - show index disk usage
copy - copy (reindex) documents from one index to another
alias - manage index aliases alias - manage index aliases
create - create an index alias create - create an index alias
list - list index aliases list - list index aliases
@@ -555,11 +582,13 @@ node - manage nodes
list - list nodes list - list nodes
show - show details about a node show - show details about a node
clients - show node http clients clients - show node http clients
usage - show node usage stats
role - manage roles role - manage roles
list - list roles list - list roles
show - show details about a role show - show details about a role
diff - show differences between roles and CSV baseline diff - show differences between roles and CSV baseline
search - search within an index search - search within an index
searchql - search using ES/QL language
shard - manage shards shard - manage shards
list - list shards list - list shards
show - show details about a shard show - show details about a shard
@@ -574,6 +603,7 @@ version - show esctl version information
debug - developer only debug - developer only
help-jsonpath - show jsonpath help help-jsonpath - show jsonpath help
help-usage - show overview of all available commands help-usage - show overview of all available commands
help-esql - show esql help
``` ```
# Development # Development

View File

@@ -5,5 +5,3 @@
- add datastream support: - add datastream support:
https://www.elastic.co/docs/api/doc/elasticsearch/operation/operation-indices-get-data-stream https://www.elastic.co/docs/api/doc/elasticsearch/operation/operation-indices-get-data-stream
also exclude data stream backing indices from index ls

View File

@@ -24,10 +24,11 @@ import (
"github.com/go-openapi/swag/loading" "github.com/go-openapi/swag/loading"
) )
//go:embed *.json //go:embed *.json *.md
var AssetFS embed.FS var AssetFS embed.FS
var OpenAPI *loads.Document var OpenAPI *loads.Document
var EsQlCheatSheet string
func LoadAssetOpenApi() { func LoadAssetOpenApi() {
doc, err := loads.Spec( doc, err := loads.Spec(
@@ -41,3 +42,12 @@ func LoadAssetOpenApi() {
OpenAPI = doc OpenAPI = doc
} }
func LoadEsql() {
md, err := AssetFS.ReadFile("esql.md")
if err != nil {
panic(err)
}
EsQlCheatSheet = string(md)
}

951
assets/esql.md Normal file
View File

@@ -0,0 +1,951 @@
<!-- courtesy https://github.com/linkan-per/ES-QL-Cheat-Sheet by |linkan-per -->
# ES|QL (Elasticsearch Query Language) Cheat Sheet
## Table of Contents
1. [Introduction](#introduction)
2. [Basic Query Structure](#basic-query-structure)
3. [Source Commands](#source-commands)
4. [Processing Commands](#processing-commands)
5. [Data Selection & Filtering](#data-selection--filtering)
6. [Aggregations & Statistics](#aggregations--statistics)
7. [String Functions](#string-functions)
8. [Mathematical Functions](#mathematical-functions)
9. [Date/Time Functions](#datetime-functions)
10. [Type Conversion Functions](#type-conversion-functions)
11. [Conditional Functions](#conditional-functions)
12. [Array Functions](#array-functions)
13. [Join Operations (LOOKUP)](#join-operations-lookup)
14. [Sorting & Limiting](#sorting--limiting)
15. [Grouping & Aggregating](#grouping--aggregating)
16. [Advanced Patterns](#advanced-patterns)
---
## Introduction
ES|QL is Elasticsearch's new query language designed for data exploration, analysis, and transformation. It uses a pipe (`|`) syntax to chain commands together.
**Basic Syntax:**
```
FROM <data-source>
| <processing-command>
| <processing-command>
| ...
```
---
## Basic Query Structure
### Simple Query
```esql
FROM logs-*
| LIMIT 10
```
### Query with Multiple Commands
```esql
FROM employees
| WHERE department == "Engineering"
| KEEP name, salary, hire_date
| SORT salary DESC
| LIMIT 5
```
---
## Source Commands
### FROM - Specify Data Source
```esql
// From a single index
FROM logs-2024
// From multiple indices with wildcard
FROM logs-*, metrics-*
// From specific indices
FROM index1, index2, index3
// With metadata
FROM logs-* METADATA _id, _index
```
### ROW - Generate Inline Data
```esql
// Create a single row
ROW name = "John", age = 30, city = "NYC"
// Multiple rows
ROW a = 1, b = "x"
| EVAL c = a * 10
```
---
## Processing Commands
### KEEP - Select Specific Fields
```esql
FROM employees
| KEEP name, department, salary
// Keep with pattern
FROM logs-*
| KEEP @timestamp, message, host.*
```
### DROP - Remove Specific Fields
```esql
FROM employees
| DROP password, ssn, internal_notes
// Drop with pattern
FROM logs-*
| DROP *.keyword
```
### RENAME - Rename Fields
```esql
FROM employees
| RENAME emp_name AS name, emp_dept AS department
// Multiple renames
FROM logs-*
| RENAME source.ip AS src_ip, destination.ip AS dst_ip
```
---
## Data Selection & Filtering
### WHERE - Filter Rows
```esql
// Equality
FROM employees
| WHERE department == "Sales"
// Comparison operators
FROM products
| WHERE price > 100 AND stock < 50
// IS NULL / IS NOT NULL
FROM logs-*
| WHERE error_code IS NOT NULL
// IN operator
FROM employees
| WHERE department IN ("Sales", "Marketing", "HR")
// LIKE operator (wildcards)
FROM logs-*
| WHERE message LIKE "*error*"
// RLIKE operator (regex)
FROM logs-*
| WHERE message RLIKE "error|exception|failure"
// NOT operator
FROM employees
| WHERE NOT department == "IT"
// Multiple conditions
FROM orders
| WHERE status == "completed"
AND total_amount > 1000
AND order_date >= "2024-01-01"
```
---
## Aggregations & Statistics
### STATS - Aggregate Functions
#### COUNT
```esql
// Count all rows
FROM logs-*
| STATS count = COUNT()
// Count distinct
FROM employees
| STATS unique_departments = COUNT_DISTINCT(department)
// Count by group
FROM logs-*
| STATS event_count = COUNT() BY log_level
```
#### SUM, AVG, MIN, MAX
```esql
FROM sales
| STATS
total_revenue = SUM(amount),
avg_sale = AVG(amount),
min_sale = MIN(amount),
max_sale = MAX(amount)
// With grouping
FROM sales
| STATS
total = SUM(amount),
average = AVG(amount)
BY product_category
```
#### MEDIAN, PERCENTILE
```esql
FROM response_times
| STATS
median_time = MEDIAN(duration),
p95 = PERCENTILE(duration, 95),
p99 = PERCENTILE(duration, 99)
```
#### Multiple Aggregations
```esql
FROM orders
| STATS
order_count = COUNT(),
total_revenue = SUM(amount),
avg_order_value = AVG(amount),
unique_customers = COUNT_DISTINCT(customer_id)
BY region, product_category
```
---
## String Functions
### CONCAT - Concatenate Strings
```esql
FROM employees
| EVAL full_name = CONCAT(first_name, " ", last_name)
// With separator
FROM logs-*
| EVAL log_info = CONCAT(level, ": ", message)
```
### SUBSTRING - Extract Substring
```esql
FROM employees
| EVAL first_initial = SUBSTRING(first_name, 0, 1)
// Extract with length
FROM products
| EVAL short_code = SUBSTRING(product_id, 0, 5)
```
### LENGTH - String Length
```esql
FROM messages
| EVAL message_length = LENGTH(message)
| WHERE message_length > 100
```
### TRIM, LTRIM, RTRIM - Remove Whitespace
```esql
FROM user_input
| EVAL cleaned = TRIM(input_field)
| EVAL left_trimmed = LTRIM(input_field)
| EVAL right_trimmed = RTRIM(input_field)
```
### UPPER, LOWER - Case Conversion
```esql
FROM employees
| EVAL name_upper = UPPER(name)
| EVAL email_lower = LOWER(email)
```
### REPLACE - Replace String
```esql
FROM logs-*
| EVAL cleaned_message = REPLACE(message, "ERROR", "Warning")
```
### SPLIT - Split String into Array
```esql
FROM logs-*
| EVAL tags_array = SPLIT(tags, ",")
```
### STARTS_WITH, ENDS_WITH
```esql
FROM files
| WHERE STARTS_WITH(filename, "log_")
| WHERE ENDS_WITH(filename, ".txt")
```
---
## Mathematical Functions
### Basic Operations
```esql
FROM sales
| EVAL
total = price * quantity,
discount_price = price * 0.9,
tax = price * 0.08
// Multiple operations
FROM metrics
| EVAL
sum_val = field1 + field2,
diff_val = field1 - field2,
product = field1 * field2,
ratio = field1 / field2,
remainder = field1 % field2
```
### ABS - Absolute Value
```esql
FROM transactions
| EVAL abs_amount = ABS(transaction_amount)
```
### ROUND, FLOOR, CEIL
```esql
FROM measurements
| EVAL
rounded = ROUND(value, 2),
floored = FLOOR(value),
ceiled = CEIL(value)
```
### POW - Power
```esql
FROM data
| EVAL squared = POW(value, 2)
| EVAL cubed = POW(value, 3)
```
### SQRT - Square Root
```esql
FROM measurements
| EVAL sqrt_value = SQRT(value)
```
### LOG, LOG10
```esql
FROM data
| EVAL
natural_log = LOG(value),
log_base_10 = LOG10(value)
```
### GREATEST, LEAST
```esql
FROM comparisons
| EVAL
max_val = GREATEST(val1, val2, val3),
min_val = LEAST(val1, val2, val3)
```
---
## Date/Time Functions
### NOW - Current Timestamp
```esql
FROM logs-*
| EVAL current_time = NOW()
```
### DATE_EXTRACT - Extract Date Parts
```esql
FROM events
| EVAL
year = DATE_EXTRACT("year", @timestamp),
month = DATE_EXTRACT("month", @timestamp),
day = DATE_EXTRACT("day", @timestamp),
hour = DATE_EXTRACT("hour", @timestamp),
minute = DATE_EXTRACT("minute", @timestamp),
day_of_week = DATE_EXTRACT("day_of_week", @timestamp)
```
### DATE_FORMAT - Format Date
```esql
FROM events
| EVAL formatted_date = DATE_FORMAT("yyyy-MM-dd", @timestamp)
| EVAL custom_format = DATE_FORMAT("MMM dd, yyyy HH:mm", @timestamp)
```
### DATE_TRUNC - Truncate Date
```esql
FROM logs-*
| EVAL
hour_bucket = DATE_TRUNC("hour", @timestamp),
day_bucket = DATE_TRUNC("day", @timestamp),
month_bucket = DATE_TRUNC("month", @timestamp)
```
### DATE_DIFF - Date Difference
```esql
FROM orders
| EVAL days_since_order = DATE_DIFF("days", order_date, NOW())
| EVAL hours_to_delivery = DATE_DIFF("hours", order_date, delivery_date)
```
### DATE_PARSE - Parse String to Date
```esql
FROM data
| EVAL parsed_date = DATE_PARSE("yyyy-MM-dd", date_string)
```
---
## Type Conversion Functions
### TO_STRING - Convert to String
```esql
FROM data
| EVAL id_string = TO_STRING(id)
| EVAL amount_string = TO_STRING(amount)
```
### TO_INTEGER, TO_LONG - Convert to Integer
```esql
FROM data
| EVAL age_int = TO_INTEGER(age_string)
| EVAL id_long = TO_LONG(id_string)
```
### TO_DOUBLE - Convert to Double
```esql
FROM data
| EVAL price_double = TO_DOUBLE(price_string)
```
### TO_BOOLEAN - Convert to Boolean
```esql
FROM data
| EVAL is_active = TO_BOOLEAN(active_string)
```
### TO_DATETIME - Convert to DateTime
```esql
FROM data
| EVAL timestamp = TO_DATETIME(date_string)
```
### TO_IP - Convert to IP Address
```esql
FROM logs-*
| EVAL ip_address = TO_IP(ip_string)
```
---
## Conditional Functions
### CASE - Conditional Logic
```esql
FROM employees
| EVAL salary_grade = CASE(
salary < 50000, "Entry",
salary < 80000, "Mid",
salary < 120000, "Senior",
"Executive"
)
// With multiple conditions
FROM orders
| EVAL order_status = CASE(
status == "pending" AND days_old > 7, "Overdue",
status == "pending", "Processing",
status == "shipped", "In Transit",
status == "delivered", "Completed",
"Unknown"
)
```
### COALESCE - Return First Non-Null Value
```esql
FROM data
| EVAL display_name = COALESCE(nickname, first_name, username, "Unknown")
```
### IF - Simple Conditional
```esql
FROM products
| EVAL stock_status =
CASE(stock > 0, "Available", "Out of Stock")
// Nested conditions
FROM employees
| EVAL bonus = CASE(
performance_rating >= 4.5, salary * 0.15,
performance_rating >= 3.5, salary * 0.10,
performance_rating >= 2.5, salary * 0.05,
0
)
```
---
## Array Functions
### MV_COUNT - Count Array Elements
```esql
FROM logs-*
| EVAL tag_count = MV_COUNT(tags)
| WHERE tag_count > 3
```
### MV_AVG, MV_SUM, MV_MIN, MV_MAX - Array Aggregations
```esql
FROM metrics
| EVAL
avg_value = MV_AVG(values),
total = MV_SUM(values),
min_value = MV_MIN(values),
max_value = MV_MAX(values)
```
### MV_CONCAT - Concatenate Array Elements
```esql
FROM logs-*
| EVAL all_tags = MV_CONCAT(tags, ", ")
```
### MV_DEDUPE - Remove Duplicates from Array
```esql
FROM data
| EVAL unique_values = MV_DEDUPE(values)
```
### MV_FIRST, MV_LAST - Get First/Last Element
```esql
FROM logs-*
| EVAL first_tag = MV_FIRST(tags)
| EVAL last_tag = MV_LAST(tags)
```
### MV_SLICE - Extract Array Slice
```esql
FROM data
| EVAL first_three = MV_SLICE(values, 0, 3)
```
---
## Join Operations (LOOKUP)
ES|QL uses ENRICH (similar to LOOKUP/JOIN) to join data from enrich policies.
### Prerequisites: Create Enrich Policy
First, create an enrich policy in Kibana Dev Tools:
```json
PUT /_enrich/policy/user_lookup
{
"match": {
"indices": "users",
"match_field": "user_id",
"enrich_fields": ["username", "email", "department"]
}
}
POST /_enrich/policy/user_lookup/_execute
```
### ENRICH - Join/Lookup Data
```esql
FROM logs-*
| ENRICH user_lookup ON user_id
| KEEP @timestamp, user_id, username, email, message
// With field renaming
FROM transactions
| ENRICH product_lookup ON product_id WITH product_name, category, price
| KEEP transaction_id, product_name, category, quantity, price
// Multiple enrichments
FROM orders
| ENRICH customer_lookup ON customer_id WITH customer_name, customer_tier
| ENRICH product_lookup ON product_id WITH product_name, product_category
| KEEP order_id, customer_name, product_name, order_amount
```
### Complex Join Example
```esql
FROM orders
| ENRICH customer_lookup ON customer_id
WITH customer_name, customer_email, customer_segment
| ENRICH product_lookup ON product_id
WITH product_name, product_category, product_price
| EVAL total_price = quantity * product_price
| WHERE customer_segment == "Premium"
| STATS
total_orders = COUNT(),
total_revenue = SUM(total_price)
BY customer_name, product_category
| SORT total_revenue DESC
```
---
## Sorting & Limiting
### SORT - Order Results
```esql
// Ascending order (default)
FROM employees
| SORT salary
// Descending order
FROM employees
| SORT salary DESC
// Multiple fields
FROM employees
| SORT department ASC, salary DESC
// With nulls first/last
FROM data
| SORT value DESC NULLS FIRST
```
### LIMIT - Limit Results
```esql
// Get first 10 rows
FROM logs-*
| LIMIT 10
// Top 5 highest salaries
FROM employees
| SORT salary DESC
| LIMIT 5
// Pagination (skip and limit)
FROM products
| SORT price
| LIMIT 20 // Results 0-19
```
### HEAD - Get First N Rows (Alias for LIMIT)
```esql
FROM logs-*
| HEAD 100
```
---
## Grouping & Aggregating
### GROUP BY with STATS
```esql
// Single field grouping
FROM sales
| STATS total_sales = SUM(amount) BY region
// Multiple field grouping
FROM orders
| STATS
order_count = COUNT(),
total_revenue = SUM(amount)
BY region, product_category, sales_rep
// Time-based grouping
FROM logs-*
| EVAL hour = DATE_TRUNC("hour", @timestamp)
| STATS event_count = COUNT() BY hour, log_level
| SORT hour DESC
```
### Complex Aggregation Example
```esql
FROM sales_data
| WHERE order_date >= "2024-01-01"
| EVAL month = DATE_TRUNC("month", order_date)
| STATS
total_orders = COUNT(),
total_revenue = SUM(amount),
avg_order_value = AVG(amount),
unique_customers = COUNT_DISTINCT(customer_id),
max_order = MAX(amount),
min_order = MIN(amount)
BY month, region, product_category
| EVAL revenue_per_customer = total_revenue / unique_customers
| WHERE total_orders > 100
| SORT month DESC, total_revenue DESC
| LIMIT 50
```
---
## Advanced Patterns
### Window Functions Pattern
```esql
// Running total by group
FROM sales
| SORT date
| STATS
daily_sales = SUM(amount),
order_count = COUNT()
BY date, region
| SORT region, date
```
### Pivoting Data
```esql
// Count by status and priority
FROM tickets
| STATS ticket_count = COUNT() BY status, priority
| SORT status, priority
```
### Finding Duplicates
```esql
FROM users
| STATS count = COUNT() BY email
| WHERE count > 1
| SORT count DESC
```
### Time Series Analysis
```esql
FROM metrics-*
| EVAL
hour = DATE_TRUNC("hour", @timestamp),
day = DATE_EXTRACT("day", @timestamp)
| STATS
avg_cpu = AVG(cpu_percent),
max_cpu = MAX(cpu_percent),
avg_memory = AVG(memory_percent)
BY hour, host
| WHERE avg_cpu > 80
| SORT hour DESC
```
### Percentage Calculations
```esql
FROM sales
| STATS
total_sales = SUM(amount),
count = COUNT()
BY product_category
| EVAL percentage = ROUND(total_sales / SUM(total_sales) * 100, 2)
| SORT percentage DESC
```
### Top N per Group
```esql
// Top 3 products per category by sales
FROM sales
| STATS total_sales = SUM(amount) BY product_category, product_name
| SORT product_category, total_sales DESC
// Note: ES|QL doesn't have native PARTITION BY,
// so you may need to process this in multiple queries or use aggregations
```
### Data Cleaning
```esql
FROM raw_data
| EVAL
// Clean whitespace
cleaned_name = TRIM(name),
// Standardize case
email_lower = LOWER(email),
// Replace values
status = REPLACE(status, "N/A", "Unknown"),
// Handle nulls
age = COALESCE(age, 0),
// Validate ranges
valid_age = CASE(age < 0 OR age > 150, NULL, age)
| WHERE cleaned_name IS NOT NULL
| DROP name, email
| RENAME cleaned_name AS name, email_lower AS email
```
### Cohort Analysis
```esql
FROM user_events
| EVAL
signup_month = DATE_TRUNC("month", signup_date),
event_month = DATE_TRUNC("month", event_date)
| STATS
active_users = COUNT_DISTINCT(user_id)
BY signup_month, event_month
| SORT signup_month, event_month
```
### Anomaly Detection Pattern
```esql
FROM metrics-*
| EVAL hour = DATE_TRUNC("hour", @timestamp)
| STATS
avg_value = AVG(value),
stddev = SQRT(AVG(POW(value - AVG(value), 2)))
BY hour
| EVAL
upper_bound = avg_value + (2 * stddev),
lower_bound = avg_value - (2 * stddev)
```
---
## Complete Real-World Examples
### Example 1: User Activity Dashboard
```esql
FROM user_logs-*
| WHERE @timestamp >= NOW() - 7 days
| ENRICH user_lookup ON user_id WITH username, user_tier
| EVAL day = DATE_TRUNC("day", @timestamp)
| STATS
daily_active_users = COUNT_DISTINCT(user_id),
total_sessions = COUNT(),
avg_session_duration = AVG(session_duration)
BY day, user_tier
| EVAL avg_duration_minutes = ROUND(avg_session_duration / 60, 2)
| SORT day DESC, user_tier
```
### Example 2: E-commerce Sales Report
```esql
FROM orders
| WHERE order_date >= "2024-01-01"
| ENRICH customer_lookup ON customer_id
WITH customer_name, customer_segment
| ENRICH product_lookup ON product_id
WITH product_name, product_category, cost_price
| EVAL
profit = (price - cost_price) * quantity,
month = DATE_TRUNC("month", order_date)
| STATS
total_orders = COUNT(),
total_revenue = SUM(price * quantity),
total_profit = SUM(profit),
avg_order_value = AVG(price * quantity),
unique_customers = COUNT_DISTINCT(customer_id)
BY month, product_category, customer_segment
| EVAL profit_margin = ROUND(total_profit / total_revenue * 100, 2)
| WHERE total_revenue > 10000
| SORT month DESC, total_revenue DESC
| LIMIT 100
```
### Example 3: Security Log Analysis
```esql
FROM security-logs-*
| WHERE @timestamp >= NOW() - 24 hours
| WHERE event_type IN ("login_failed", "suspicious_activity")
| EVAL hour = DATE_TRUNC("hour", @timestamp)
| STATS
event_count = COUNT(),
unique_ips = COUNT_DISTINCT(source_ip),
unique_users = COUNT_DISTINCT(username)
BY hour, event_type, country
| WHERE event_count > 100
| SORT hour DESC, event_count DESC
```
### Example 4: Application Performance Monitoring
```esql
FROM apm-*
| WHERE @timestamp >= NOW() - 1 hour
| EVAL
response_category = CASE(
response_time < 100, "Fast",
response_time < 500, "Medium",
response_time < 1000, "Slow",
"Very Slow"
),
minute = DATE_TRUNC("minute", @timestamp)
| STATS
request_count = COUNT(),
avg_response = AVG(response_time),
p95_response = PERCENTILE(response_time, 95),
p99_response = PERCENTILE(response_time, 99),
error_count = COUNT() WHERE status_code >= 400
BY minute, endpoint, response_category
| EVAL error_rate = ROUND(error_count / request_count * 100, 2)
| WHERE request_count > 10
| SORT minute DESC, avg_response DESC
```
---
## Tips & Best Practices
1. **Use KEEP instead of SELECT** - More explicit about which fields to retain
2. **Filter early with WHERE** - Reduce data processing by filtering before aggregations
3. **Use DATE_TRUNC for time bucketing** - Essential for time series analysis
4. **Leverage ENRICH for joins** - Pre-create enrich policies for frequently joined data
5. **Use EVAL for calculated fields** - Create derived fields before aggregation
6. **Combine multiple conditions in WHERE** - More efficient than multiple WHERE clauses
7. **Use STATS with BY for grouping** - Replaces traditional GROUP BY
8. **Sort after aggregation** - More efficient than sorting before
9. **Use LIMIT to control output size** - Especially important for large datasets
10. **Use metadata fields when needed** - Access _id, _index with METADATA keyword
---
## Common Patterns Cheat Sheet
```esql
// Count by field
FROM index | STATS count = COUNT() BY field
// Top N
FROM index | STATS value = SUM(amount) BY category | SORT value DESC | LIMIT 10
// Time series
FROM index | EVAL bucket = DATE_TRUNC("hour", @timestamp) | STATS count = COUNT() BY bucket
// Percentage of total
FROM index | STATS total = SUM(amount) BY category | EVAL pct = total / SUM(total) * 100
// Filter nulls
FROM index | WHERE field IS NOT NULL
// String matching
FROM index | WHERE field LIKE "*pattern*"
// Date range
FROM index | WHERE @timestamp >= NOW() - 7 days
// Multiple aggregations
FROM index | STATS count = COUNT(), sum = SUM(val), avg = AVG(val) BY group
// Conditional aggregation
FROM index | STATS error_count = COUNT() WHERE status == "error" BY service
```
---
## Comparison with Traditional SQL
| SQL | ES|QL |
|-----|-------|
| SELECT * | FROM index |
| SELECT field1, field2 | FROM index \| KEEP field1, field2 |
| WHERE condition | WHERE condition (same) |
| GROUP BY field | STATS ... BY field |
| ORDER BY field | SORT field |
| LIMIT 10 | LIMIT 10 (same) |
| COUNT(*) | STATS count = COUNT() |
| SUM(field) | STATS total = SUM(field) |
| AVG(field) | STATS avg = AVG(field) |
| JOIN | ENRICH (using enrich policies) |
| CASE WHEN | CASE(...) |
| CONCAT(a, b) | CONCAT(a, b) (same) |
---
## Resources
- Official ES|QL Documentation: https://www.elastic.co/guide/en/elasticsearch/reference/current/esql.html
- ES|QL Functions Reference: https://www.elastic.co/guide/en/elasticsearch/reference/current/esql-functions.html
- Enrich Processor: https://www.elastic.co/guide/en/elasticsearch/reference/current/enrich-processor.html
---
*Last Updated: 2024*
*ES|QL is actively evolving - check official documentation for latest features*

View File

@@ -72,7 +72,7 @@ func ApiShow(conf *cfg.Config) *cli.Command {
}, },
ShellComplete: func(ctx context.Context, cmd *cli.Command) { ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Capi) complete(conf, cmd, Capi)
}, },
} }
} }
@@ -91,6 +91,12 @@ func ApiRepl(conf *cfg.Config) *cli.Command {
Aliases: []string{"p"}, Aliases: []string{"p"},
Sources: cli.EnvVars("PAGER", "ES_JSON_PAGER"), Sources: cli.EnvVars("PAGER", "ES_JSON_PAGER"),
}, },
&cli.BoolFlag{
Name: "human-readable-cat",
Usage: "enable human readable /_cat output",
Destination: &conf.HumanCat,
Aliases: []string{"H"},
},
}, },
Action: func(ctx context.Context, cmd *cli.Command) error { Action: func(ctx context.Context, cmd *cli.Command) error {

View File

@@ -52,14 +52,14 @@ func CcrStatus(conf *cfg.Config) *cli.Command {
Flags: []cli.Flag{ Flags: []cli.Flag{
&cli.StringFlag{ &cli.StringFlag{
Name: "exclude", Name: "exclude",
Usage: "regexp of indicies to exclude", Usage: "regexp of indices to exclude",
Destination: &conf.Exclude, Destination: &conf.Exclude,
Aliases: []string{"e"}, Aliases: []string{"e"},
}, },
}, },
ShellComplete: func(ctx context.Context, cmd *cli.Command) { ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Ccluster) complete(conf, cmd, Ccluster)
}, },
Action: func(ctx context.Context, cmd *cli.Command) error { Action: func(ctx context.Context, cmd *cli.Command) error {
@@ -115,7 +115,7 @@ func CcrRemoteInfo(conf *cfg.Config) *cli.Command {
UsageText: "info [options] [<index>]", UsageText: "info [options] [<index>]",
ShellComplete: func(ctx context.Context, cmd *cli.Command) { ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cindex) complete(conf, cmd, Cindex)
}, },
Action: func(ctx context.Context, cmd *cli.Command) error { Action: func(ctx context.Context, cmd *cli.Command) error {

View File

@@ -60,7 +60,7 @@ func CcrFollowerRenew(conf *cfg.Config) *cli.Command {
}, },
ShellComplete: func(ctx context.Context, cmd *cli.Command) { ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cindex) complete(conf, cmd, Cindex)
}, },
Action: func(ctx context.Context, cmd *cli.Command) error { Action: func(ctx context.Context, cmd *cli.Command) error {
@@ -82,7 +82,7 @@ func CcrFollowerResume(conf *cfg.Config) *cli.Command {
UsageText: "resume [options] <index>", UsageText: "resume [options] <index>",
ShellComplete: func(ctx context.Context, cmd *cli.Command) { ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cindex) complete(conf, cmd, Cindex)
}, },
Action: func(ctx context.Context, cmd *cli.Command) error { Action: func(ctx context.Context, cmd *cli.Command) error {
@@ -104,7 +104,7 @@ func CcrFollowerPause(conf *cfg.Config) *cli.Command {
UsageText: "pause [options] <index>", UsageText: "pause [options] <index>",
ShellComplete: func(ctx context.Context, cmd *cli.Command) { ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cindex) complete(conf, cmd, Cindex)
}, },
Action: func(ctx context.Context, cmd *cli.Command) error { Action: func(ctx context.Context, cmd *cli.Command) error {
@@ -126,7 +126,7 @@ func CcrFollowerUnfollow(conf *cfg.Config) *cli.Command {
UsageText: "unfollow [options] <index>", UsageText: "unfollow [options] <index>",
ShellComplete: func(ctx context.Context, cmd *cli.Command) { ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cindex) complete(conf, cmd, Cindex)
}, },
Action: func(ctx context.Context, cmd *cli.Command) error { Action: func(ctx context.Context, cmd *cli.Command) error {
@@ -158,7 +158,7 @@ func CcrFollowerAdd(conf *cfg.Config) *cli.Command {
}, },
ShellComplete: func(ctx context.Context, cmd *cli.Command) { ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cindex) complete(conf, cmd, Cindex)
}, },
Action: func(ctx context.Context, cmd *cli.Command) error { Action: func(ctx context.Context, cmd *cli.Command) error {
@@ -181,7 +181,7 @@ func CcrFollowerDelete(conf *cfg.Config) *cli.Command {
UsageText: "delete <index>", UsageText: "delete <index>",
ShellComplete: func(ctx context.Context, cmd *cli.Command) { ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cindex) complete(conf, cmd, Cindex)
}, },
Action: func(ctx context.Context, cmd *cli.Command) error { Action: func(ctx context.Context, cmd *cli.Command) error {
@@ -205,7 +205,7 @@ func CcrFollowerShow(conf *cfg.Config) *cli.Command {
UsageText: "show <index>", UsageText: "show <index>",
ShellComplete: func(ctx context.Context, cmd *cli.Command) { ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cindex) complete(conf, cmd, Cindex)
}, },
Action: func(ctx context.Context, cmd *cli.Command) error { Action: func(ctx context.Context, cmd *cli.Command) error {

View File

@@ -83,7 +83,7 @@ func ClusterSwitch(conf *cfg.Config) *cli.Command {
Aliases: []string{"ctx"}, Aliases: []string{"ctx"},
ShellComplete: func(ctx context.Context, cmd *cli.Command) { ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Ccluster) complete(conf, cmd, Ccluster)
}, },
Action: func(ctx context.Context, cmd *cli.Command) error { Action: func(ctx context.Context, cmd *cli.Command) error {

View File

@@ -50,7 +50,7 @@ func ClusterRerouteMove(conf *cfg.Config) *cli.Command {
UsageText: "move [options] <index>", UsageText: "move [options] <index>",
ShellComplete: func(ctx context.Context, cmd *cli.Command) { ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cindex) complete(conf, cmd, Cindex)
}, },
Flags: []cli.Flag{ Flags: []cli.Flag{
@@ -95,7 +95,7 @@ func ClusterRerouteAllocateReplica(conf *cfg.Config) *cli.Command {
UsageText: "allocate-replica [options] <index>", UsageText: "allocate-replica [options] <index>",
ShellComplete: func(ctx context.Context, cmd *cli.Command) { ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cindex) complete(conf, cmd, Cindex)
}, },
Flags: []cli.Flag{ Flags: []cli.Flag{
@@ -133,7 +133,7 @@ func ClusterRerouteCancel(conf *cfg.Config) *cli.Command {
UsageText: "cancel [options] <index>", UsageText: "cancel [options] <index>",
ShellComplete: func(ctx context.Context, cmd *cli.Command) { ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cindex) complete(conf, cmd, Cindex)
}, },
Flags: []cli.Flag{ Flags: []cli.Flag{
@@ -185,7 +185,7 @@ func ClusterRerouteAllocatePrimary(conf *cfg.Config, stale bool) *cli.Command {
UsageText: name + " [options] <index>", UsageText: name + " [options] <index>",
ShellComplete: func(ctx context.Context, cmd *cli.Command) { ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cindex) complete(conf, cmd, Cindex)
}, },
Flags: []cli.Flag{ Flags: []cli.Flag{

View File

@@ -105,7 +105,7 @@ func ClusterSettingsSet(conf *cfg.Config) *cli.Command {
}, },
ShellComplete: func(ctx context.Context, cmd *cli.Command) { ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cclustersettings) complete(conf, cmd, Cclustersettings)
}, },
Action: func(ctx context.Context, cmd *cli.Command) error { Action: func(ctx context.Context, cmd *cli.Command) error {

View File

@@ -34,20 +34,14 @@ const (
Cilm Cilm
Cnode Cnode
Cclustersettings Cclustersettings
Cindextemplate
) )
func complete(cmd *cli.Command, what int) { func complete(conf *cfg.Config, cmd *cli.Command, what int) {
if cmd.NArg() > 0 { if cmd.NArg() > 0 {
return return
} }
// FIXME: config should load from root.Before(), see https://github.com/urfave/cli/issues/2348
// workaround: load it directly here
conf := cfg.NewConfig()
if err := conf.Init(); err != nil {
return
}
var ( var (
list []string list []string
err error err error
@@ -72,6 +66,8 @@ func complete(cmd *cli.Command, what int) {
list, err = es.NodeNames(conf) list, err = es.NodeNames(conf)
case Cclustersettings: case Cclustersettings:
list, err = es.ClusterSettingsNames(conf) list, err = es.ClusterSettingsNames(conf)
case Cindextemplate:
list, err = es.IndexTemplateNames(conf)
} }
if err != nil { if err != nil {

View File

@@ -46,7 +46,7 @@ func DatastreamList(conf *cfg.Config) *cli.Command {
return &cli.Command{ return &cli.Command{
Name: "list", Name: "list",
Aliases: []string{"ls"}, Aliases: []string{"ls"},
Usage: "list indicies", Usage: "list indices",
Flags: []cli.Flag{ Flags: []cli.Flag{
&cli.IntFlag{ &cli.IntFlag{
@@ -91,7 +91,7 @@ func DatastreamShow(conf *cfg.Config) *cli.Command {
}, },
ShellComplete: func(ctx context.Context, cmd *cli.Command) { ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cdatastream) complete(conf, cmd, Cdatastream)
}, },
} }
} }
@@ -113,7 +113,7 @@ func DatastreamCreate(conf *cfg.Config) *cli.Command {
}, },
ShellComplete: func(ctx context.Context, cmd *cli.Command) { ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cdatastream) complete(conf, cmd, Cdatastream)
}, },
} }
} }
@@ -135,7 +135,7 @@ func DatastreamDelete(conf *cfg.Config) *cli.Command {
}, },
ShellComplete: func(ctx context.Context, cmd *cli.Command) { ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cdatastream) complete(conf, cmd, Cdatastream)
}, },
} }
} }
@@ -148,7 +148,7 @@ func DatastreamRollover(conf *cfg.Config) *cli.Command {
UsageText: "rollover <data stream>", UsageText: "rollover <data stream>",
ShellComplete: func(ctx context.Context, cmd *cli.Command) { ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cindex) complete(conf, cmd, Cindex)
}, },
Flags: []cli.Flag{ Flags: []cli.Flag{
@@ -169,7 +169,7 @@ func DatastreamRollover(conf *cfg.Config) *cli.Command {
Usage: "roll over after max age (eg: 7d, 2m, 8h)", Usage: "roll over after max age (eg: 7d, 2m, 8h)",
Destination: &conf.MaxAge, Destination: &conf.MaxAge,
}, },
&cli.IntFlag{ &cli.Int64Flag{
Name: "max-docs", Name: "max-docs",
Usage: "roll over after max docs", Usage: "roll over after max docs",
Destination: &conf.MaxDocs, Destination: &conf.MaxDocs,
@@ -205,7 +205,7 @@ func DatastreamIlm(conf *cfg.Config) *cli.Command {
UsageText: "ds ilm <name>", UsageText: "ds ilm <name>",
ShellComplete: func(ctx context.Context, cmd *cli.Command) { ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cdatastream) complete(conf, cmd, Cdatastream)
}, },
Action: func(ctx context.Context, cmd *cli.Command) error { Action: func(ctx context.Context, cmd *cli.Command) error {

View File

@@ -123,7 +123,7 @@ func IlmShow(conf *cfg.Config) *cli.Command {
}, },
ShellComplete: func(ctx context.Context, cmd *cli.Command) { ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cilm) complete(conf, cmd, Cilm)
}, },
} }
} }

View File

@@ -89,7 +89,7 @@ func IlmForecastList(conf *cfg.Config) *cli.Command {
}, },
&cli.BoolFlag{ &cli.BoolFlag{
Name: "hidden", Name: "hidden",
Usage: "include hidden indicies", Usage: "include hidden indices",
Destination: &conf.Hidden, Destination: &conf.Hidden,
Aliases: []string{"H"}, Aliases: []string{"H"},
}, },

View File

@@ -30,7 +30,7 @@ func Index(conf *cfg.Config) *cli.Command {
return &cli.Command{ return &cli.Command{
Name: "index", Name: "index",
Aliases: []string{"i"}, Aliases: []string{"i"},
Usage: "manage indicies", Usage: "manage indices",
Commands: []*cli.Command{ Commands: []*cli.Command{
IndexList(conf), IndexList(conf),
@@ -41,6 +41,8 @@ func Index(conf *cfg.Config) *cli.Command {
IndexClose(conf), IndexClose(conf),
IndexFields(conf), IndexFields(conf),
IndexIlm(conf), IndexIlm(conf),
IndexDu(conf),
IndexCopy(conf),
// sub commands // sub commands
IndexAlias(conf), IndexAlias(conf),
@@ -53,7 +55,7 @@ func IndexList(conf *cfg.Config) *cli.Command {
return &cli.Command{ return &cli.Command{
Name: "list", Name: "list",
Aliases: []string{"ls"}, Aliases: []string{"ls"},
Usage: "list indicies", Usage: "list indices",
Flags: []cli.Flag{ Flags: []cli.Flag{
&cli.IntFlag{ &cli.IntFlag{
@@ -64,25 +66,25 @@ func IndexList(conf *cfg.Config) *cli.Command {
}, },
&cli.BoolFlag{ &cli.BoolFlag{
Name: "partials", Name: "partials",
Usage: "include partial indicies", Usage: "include partial indices",
Destination: &conf.Partials, Destination: &conf.Partials,
Aliases: []string{"p"}, Aliases: []string{"p"},
}, },
&cli.BoolFlag{ &cli.BoolFlag{
Name: "hidden", Name: "hidden",
Usage: "include hidden indicies", Usage: "include hidden indices",
Destination: &conf.Hidden, Destination: &conf.Hidden,
Aliases: []string{"H"}, Aliases: []string{"H"},
}, },
&cli.BoolFlag{ &cli.BoolFlag{
Name: "failed", Name: "failed",
Usage: "include only red failed indicies", Usage: "include only red failed indices",
Destination: &conf.Failed, Destination: &conf.Failed,
Aliases: []string{"f"}, Aliases: []string{"f"},
}, },
&cli.StringSliceFlag{ &cli.StringSliceFlag{
Name: "filter", Name: "filter",
Usage: "show only indicies matching the filter", Usage: "show only indices matching the filter",
Destination: &conf.Filter, Destination: &conf.Filter,
Aliases: []string{"F"}, Aliases: []string{"F"},
}, },
@@ -110,7 +112,7 @@ func IndexShow(conf *cfg.Config) *cli.Command {
}, },
ShellComplete: func(ctx context.Context, cmd *cli.Command) { ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cindex) complete(conf, cmd, Cindex)
}, },
} }
} }
@@ -180,7 +182,7 @@ func IndexDelete(conf *cfg.Config) *cli.Command {
UsageText: "delete <index>", UsageText: "delete <index>",
ShellComplete: func(ctx context.Context, cmd *cli.Command) { ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cindex) complete(conf, cmd, Cindex)
}, },
Action: func(ctx context.Context, cmd *cli.Command) error { Action: func(ctx context.Context, cmd *cli.Command) error {
@@ -201,7 +203,7 @@ func IndexClose(conf *cfg.Config) *cli.Command {
UsageText: "close <index>", UsageText: "close <index>",
ShellComplete: func(ctx context.Context, cmd *cli.Command) { ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cindex) complete(conf, cmd, Cindex)
}, },
Action: func(ctx context.Context, cmd *cli.Command) error { Action: func(ctx context.Context, cmd *cli.Command) error {
@@ -243,7 +245,7 @@ func IndexFields(conf *cfg.Config) *cli.Command {
}, },
ShellComplete: func(ctx context.Context, cmd *cli.Command) { ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cindex) complete(conf, cmd, Cindex)
}, },
Action: func(ctx context.Context, cmd *cli.Command) error { Action: func(ctx context.Context, cmd *cli.Command) error {
@@ -264,7 +266,7 @@ func IndexIlm(conf *cfg.Config) *cli.Command {
UsageText: "index ilm <index>", UsageText: "index ilm <index>",
ShellComplete: func(ctx context.Context, cmd *cli.Command) { ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cindex) complete(conf, cmd, Cindex)
}, },
Action: func(ctx context.Context, cmd *cli.Command) error { Action: func(ctx context.Context, cmd *cli.Command) error {
@@ -277,3 +279,103 @@ func IndexIlm(conf *cfg.Config) *cli.Command {
}, },
} }
} }
func IndexDu(conf *cfg.Config) *cli.Command {
return &cli.Command{
Name: "du",
Usage: "show index disk usage",
UsageText: "index du <index>",
ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(conf, cmd, Cindex)
},
Action: func(ctx context.Context, cmd *cli.Command) error {
index := cmd.Args().Get(0)
if index == "" {
return errors.New("no index specified")
}
return es.IndexDiskusage(conf, index)
},
}
}
func IndexCopy(conf *cfg.Config) *cli.Command {
return &cli.Command{
Name: "copy",
Aliases: []string{"alias", "cp"},
Usage: "copy (reindex) documents from one index to another",
UsageText: "index copy [options] -s <source-index> -t <dest-index>",
Flags: []cli.Flag{
// FIXME: not implemented by typed API
// &cli.BoolFlag{
// Name: "missing", // op_type
// Usage: "copy only missing docs",
// Destination: &conf.Missing,
// Aliases: []string{"m"},
// },
// &cli.BoolFlag{
// Name: "sync", // version_type to external
// Usage: "create missing docs and update outdated docs",
// Destination: &conf.Sync,
// Aliases: []string{"S"},
// },
&cli.BoolFlag{
Name: "force", // conflicts to proceed
Usage: "continue reindexing even when conflicts happen",
Destination: &conf.Force,
Aliases: []string{"f"},
},
&cli.Float64Flag{
Name: "requests-per-second",
Usage: "the maximum number of documents to index per second (-1 turns off throttling)",
Destination: &conf.RequestsPerSecond,
Aliases: []string{"R"},
},
&cli.Int64Flag{
Name: "max-docs",
Usage: "the maximum number of documents to reindex",
Destination: &conf.MaxDocs,
Aliases: []string{"m"},
},
&cli.DurationFlag{
Name: "timeout",
Usage: "timeout for write operations (e.g. 300m or 120s)",
Destination: &conf.Timeout,
Aliases: []string{"T"},
},
&cli.BoolFlag{
Name: "refresh",
Usage: "refresh affected shards to make this operation visible to search",
Destination: &conf.Refresh,
Aliases: []string{"r"},
},
&cli.BoolFlag{
Name: "wait",
Usage: "wait for active shards",
Destination: &conf.Wait,
Aliases: []string{"w"},
},
&cli.StringSliceFlag{
Name: "source",
Usage: "source index (multiple supported)",
Destination: &conf.SourceIndices,
Aliases: []string{"s"},
Required: true,
},
&cli.StringFlag{
Name: "target",
Usage: "target index",
Destination: &conf.Index,
Aliases: []string{"t"},
Required: true,
},
},
Action: func(ctx context.Context, cmd *cli.Command) error {
return es.IndexCopy(conf)
},
}
}

View File

@@ -52,7 +52,7 @@ func IndexAliasCreate(conf *cfg.Config) *cli.Command {
UsageText: "create <index> <alias>", UsageText: "create <index> <alias>",
ShellComplete: func(ctx context.Context, cmd *cli.Command) { ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cindex) complete(conf, cmd, Cindex)
}, },
Action: func(ctx context.Context, cmd *cli.Command) error { Action: func(ctx context.Context, cmd *cli.Command) error {
@@ -76,7 +76,7 @@ func IndexAliasDelete(conf *cfg.Config) *cli.Command {
UsageText: "delete <index> <alias>", UsageText: "delete <index> <alias>",
ShellComplete: func(ctx context.Context, cmd *cli.Command) { ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cindex) complete(conf, cmd, Cindex)
}, },
Action: func(ctx context.Context, cmd *cli.Command) error { Action: func(ctx context.Context, cmd *cli.Command) error {
@@ -100,7 +100,7 @@ func IndexAliasRollover(conf *cfg.Config) *cli.Command {
UsageText: "rollover <alias>", UsageText: "rollover <alias>",
ShellComplete: func(ctx context.Context, cmd *cli.Command) { ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cindex) complete(conf, cmd, Cindex)
}, },
Flags: []cli.Flag{ Flags: []cli.Flag{
@@ -121,7 +121,7 @@ func IndexAliasRollover(conf *cfg.Config) *cli.Command {
Usage: "roll over after max age (eg: 7d, 2m, 8h)", Usage: "roll over after max age (eg: 7d, 2m, 8h)",
Destination: &conf.MaxAge, Destination: &conf.MaxAge,
}, },
&cli.IntFlag{ &cli.Int64Flag{
Name: "max-docs", Name: "max-docs",
Usage: "roll over after max docs", Usage: "roll over after max docs",
Destination: &conf.MaxDocs, Destination: &conf.MaxDocs,
@@ -159,7 +159,7 @@ func IndexAliasList(conf *cfg.Config) *cli.Command {
Flags: []cli.Flag{ Flags: []cli.Flag{
&cli.StringSliceFlag{ &cli.StringSliceFlag{
Name: "filter", Name: "filter",
Usage: "show only aliases for indicies matching the filter", Usage: "show only aliases for indices matching the filter",
Destination: &conf.Filter, Destination: &conf.Filter,
Aliases: []string{"F"}, Aliases: []string{"F"},
}, },

View File

@@ -49,8 +49,23 @@ func IndexTemplateList(conf *cfg.Config) *cli.Command {
Aliases: []string{"ls"}, Aliases: []string{"ls"},
Usage: "list index templates", Usage: "list index templates",
Flags: []cli.Flag{
&cli.StringSliceFlag{
Name: "filter-index-pattern",
Usage: "show only index templates which use patterns, which match the filter",
Destination: &conf.Filter,
Aliases: []string{"F"},
},
&cli.BoolFlag{
Name: "hidden",
Usage: "show hidden index templates as well",
Destination: &conf.Hidden,
Aliases: []string{"H"},
},
},
Action: func(ctx context.Context, cmd *cli.Command) error { Action: func(ctx context.Context, cmd *cli.Command) error {
return es.IndexTemplateList(conf) return es.IndexTemplateList(conf, cmd.Args().Get(0))
}, },
} }
} }
@@ -71,7 +86,7 @@ func IndexTemplateShow(conf *cfg.Config) *cli.Command {
}, },
ShellComplete: func(ctx context.Context, cmd *cli.Command) { ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cindex) complete(conf, cmd, Cindextemplate)
}, },
} }
} }
@@ -215,7 +230,7 @@ func IndexTemplateDelete(conf *cfg.Config) *cli.Command {
}, },
ShellComplete: func(ctx context.Context, cmd *cli.Command) { ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cindex) complete(conf, cmd, Cindextemplate)
}, },
} }
} }

View File

@@ -36,6 +36,7 @@ func Node(conf *cfg.Config) *cli.Command {
NodeList(conf), NodeList(conf),
NodeShow(conf), NodeShow(conf),
NodeClients(conf), NodeClients(conf),
NodeUsage(conf),
}, },
} }
} }
@@ -60,7 +61,7 @@ func NodeShow(conf *cfg.Config) *cli.Command {
UsageText: "show [options] <node>", UsageText: "show [options] <node>",
ShellComplete: func(ctx context.Context, cmd *cli.Command) { ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cnode) complete(conf, cmd, Cnode)
}, },
Action: func(ctx context.Context, cmd *cli.Command) error { Action: func(ctx context.Context, cmd *cli.Command) error {
@@ -90,7 +91,7 @@ func NodeClients(conf *cfg.Config) *cli.Command {
}, },
ShellComplete: func(ctx context.Context, cmd *cli.Command) { ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Cnode) complete(conf, cmd, Cnode)
}, },
Action: func(ctx context.Context, cmd *cli.Command) error { Action: func(ctx context.Context, cmd *cli.Command) error {
@@ -103,3 +104,19 @@ func NodeClients(conf *cfg.Config) *cli.Command {
}, },
} }
} }
func NodeUsage(conf *cfg.Config) *cli.Command {
return &cli.Command{
Name: "usage",
Usage: "show node usage stats",
UsageText: "usage [options] [<node>]",
ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(conf, cmd, Cnode)
},
Action: func(ctx context.Context, cmd *cli.Command) error {
return es.NodeUsage(conf, cmd.Args().Get(0))
},
}
}

View File

@@ -68,7 +68,7 @@ func RoleShow(conf *cfg.Config) *cli.Command {
UsageText: "show [options] <role>", UsageText: "show [options] <role>",
ShellComplete: func(ctx context.Context, cmd *cli.Command) { ShellComplete: func(ctx context.Context, cmd *cli.Command) {
complete(cmd, Crole) complete(conf, cmd, Crole)
}, },
Action: func(ctx context.Context, cmd *cli.Command) error { Action: func(ctx context.Context, cmd *cli.Command) error {

View File

@@ -21,13 +21,17 @@ import (
"fmt" "fmt"
golog "log" golog "log"
"os" "os"
"runtime/debug"
"runtime/pprof" "runtime/pprof"
"strings" "strings"
"codeberg.org/scip/esctl/assets"
"codeberg.org/scip/esctl/pkg/cfg" "codeberg.org/scip/esctl/pkg/cfg"
"codeberg.org/scip/esctl/pkg/es" "codeberg.org/scip/esctl/pkg/es"
"codeberg.org/scip/esctl/pkg/log" "codeberg.org/scip/esctl/pkg/log"
"codeberg.org/scip/esctl/pkg/printer"
markdown "github.com/MichaelMure/go-term-markdown"
"github.com/urfave/cli/v3" "github.com/urfave/cli/v3"
) )
@@ -115,6 +119,12 @@ func Main() int {
Usage: "enable HTTP debugging", Usage: "enable HTTP debugging",
Destination: &conf.DebugHTTP, Destination: &conf.DebugHTTP,
}, },
&cli.BoolFlag{
Name: "debug-goroutines",
Value: false,
Usage: "enable goroutine debugging",
Destination: &conf.DebugGoRoutines,
},
&cli.BoolFlag{ &cli.BoolFlag{
Name: "align-ints", Name: "align-ints",
Aliases: []string{"I"}, Aliases: []string{"I"},
@@ -141,7 +151,7 @@ func Main() int {
Name: "output", Name: "output",
Aliases: []string{"o"}, Aliases: []string{"o"},
Value: "", Value: "",
Usage: "output mode (tsv, json, yaml) default: tsv", Usage: "output mode (tsv, csv, json, yaml) default: tsv",
Destination: &conf.Output, Destination: &conf.Output,
}, },
&cli.StringFlag{ &cli.StringFlag{
@@ -164,6 +174,7 @@ func Main() int {
Node(conf), Node(conf),
Roles(conf), Roles(conf),
Search(conf), Search(conf),
SearchQL(conf),
Shard(conf), Shard(conf),
Snapshot(conf), Snapshot(conf),
Task(conf), Task(conf),
@@ -171,6 +182,7 @@ func Main() int {
Debug(conf), Debug(conf),
HelpJsonPath(conf), HelpJsonPath(conf),
HelpUsage(conf), HelpUsage(conf),
HelpEsQl(conf),
}, },
Before: func(ctx context.Context, cmd *cli.Command) (context.Context, error) { Before: func(ctx context.Context, cmd *cli.Command) (context.Context, error) {
@@ -226,12 +238,34 @@ func HelpJsonPath(conf *cfg.Config) *cli.Command {
} }
} }
func HelpEsQl(conf *cfg.Config) *cli.Command {
return &cli.Command{
Name: "help-esql",
Usage: "show esql help",
Action: func(ctx context.Context, cmd *cli.Command) error {
assets.LoadEsql()
width := cfg.GetTermWidth()
printer.Pager("esql cheat sheet", string(markdown.Render(assets.EsQlCheatSheet, width, cfg.DefaultMargin)))
return nil
},
}
}
func Version(conf *cfg.Config) *cli.Command { func Version(conf *cfg.Config) *cli.Command {
return &cli.Command{ return &cli.Command{
Name: "version", Name: "version",
Usage: "show esctl version information", Usage: "show esctl version information",
Action: func(ctx context.Context, cmd *cli.Command) error { Action: func(ctx context.Context, cmd *cli.Command) error {
info, _ := debug.ReadBuildInfo()
if strings.Contains(info.Main.Version, "+dirty") {
cfg.COMMIT += "+dirty"
}
_, err := fmt.Printf(versionFmt, _, err := fmt.Printf(versionFmt,
cfg.Version, cfg.BUILD, cfg.BRANCH, cfg.COMMIT, cfg.GOVERSION, cfg.APIVERSION) cfg.Version, cfg.BUILD, cfg.BRANCH, cfg.COMMIT, cfg.GOVERSION, cfg.APIVERSION)

View File

@@ -18,6 +18,7 @@ package cmd
import ( import (
"context" "context"
"strings"
"codeberg.org/scip/esctl/pkg/cfg" "codeberg.org/scip/esctl/pkg/cfg"
"codeberg.org/scip/esctl/pkg/es" "codeberg.org/scip/esctl/pkg/es"
@@ -45,6 +46,11 @@ https://www.elastic.co/docs/reference/elasticsearch/rest-apis/common-options#dat
For timestamp formats refer to: For timestamp formats refer to:
https://www.elastic.co/docs/reference/elasticsearch/mapping-reference/mapping-date-format` https://www.elastic.co/docs/reference/elasticsearch/mapping-reference/mapping-date-format`
const QlUsage = `ES|QL documentation:
https://www.elastic.co/docs/reference/query-languages/esql/esql-getting-started
See also: esctl help-esql`
func Search(conf *cfg.Config) *cli.Command { func Search(conf *cfg.Config) *cli.Command {
return &cli.Command{ return &cli.Command{
Name: "search", Name: "search",
@@ -149,3 +155,19 @@ func Search(conf *cfg.Config) *cli.Command {
}, },
} }
} }
func SearchQL(conf *cfg.Config) *cli.Command {
return &cli.Command{
Name: "searchql",
Aliases: []string{"/"},
Usage: "search using ES/QL language",
UsageText: "search <ES/QL term>",
CustomHelpTemplate: addReference(QlUsage),
Action: func(ctx context.Context, cmd *cli.Command) error {
args := cmd.Args()
return es.SearchQL(conf, strings.Join(args.Slice(), " "))
},
}
}

4
go.mod
View File

@@ -17,8 +17,9 @@ module codeberg.org/scip/esctl
go 1.26 go 1.26
require ( require (
codeberg.org/scip/mapmap v0.0.2
github.com/MichaelMure/go-term-markdown v0.1.4 github.com/MichaelMure/go-term-markdown v0.1.4
github.com/alecthomas/repr v0.5.2 github.com/alecthomas/repr v0.5.3
github.com/charmbracelet/bubbles v1.0.0 github.com/charmbracelet/bubbles v1.0.0
github.com/charmbracelet/bubbletea v1.3.10 github.com/charmbracelet/bubbletea v1.3.10
github.com/charmbracelet/lipgloss v1.1.0 github.com/charmbracelet/lipgloss v1.1.0
@@ -31,7 +32,6 @@ require (
github.com/go-openapi/spec v0.22.5 github.com/go-openapi/spec v0.22.5
github.com/go-openapi/swag/loading v0.26.1 github.com/go-openapi/swag/loading v0.26.1
github.com/mattn/go-isatty v0.0.22 github.com/mattn/go-isatty v0.0.22
github.com/seeruk/go-wordwrap v0.0.0-20191208221741-14ec4aac9550
github.com/tidwall/gjson v1.19.0 github.com/tidwall/gjson v1.19.0
github.com/tlinden/yadu v0.1.3 github.com/tlinden/yadu v0.1.3
github.com/urfave/cli/v3 v3.10.1-0.20260623012112-f980ca84bf65 github.com/urfave/cli/v3 v3.10.1-0.20260623012112-f980ca84bf65

8
go.sum
View File

@@ -1,3 +1,5 @@
codeberg.org/scip/mapmap v0.0.2 h1:0i61jOUwFmGVwPukrMzrx0Fi4T9Qru2nlmibuaJimBo=
codeberg.org/scip/mapmap v0.0.2/go.mod h1:/ojYo2P7dMA2FWEu+jHKmsKPeq5yDcCvXeHevqZd5OI=
github.com/MichaelMure/go-term-markdown v0.1.4 h1:Ir3kBXDUtOX7dEv0EaQV8CNPpH+T7AfTh0eniMOtNcs= github.com/MichaelMure/go-term-markdown v0.1.4 h1:Ir3kBXDUtOX7dEv0EaQV8CNPpH+T7AfTh0eniMOtNcs=
github.com/MichaelMure/go-term-markdown v0.1.4/go.mod h1:EhcA3+pKYnlUsxYKBJ5Sn1cTQmmBMjeNlpV8nRb+JxA= github.com/MichaelMure/go-term-markdown v0.1.4/go.mod h1:EhcA3+pKYnlUsxYKBJ5Sn1cTQmmBMjeNlpV8nRb+JxA=
github.com/MichaelMure/go-term-text v0.3.1 h1:Kw9kZanyZWiCHOYu9v/8pWEgDQ6UVN9/ix2Vd2zzWf0= github.com/MichaelMure/go-term-text v0.3.1 h1:Kw9kZanyZWiCHOYu9v/8pWEgDQ6UVN9/ix2Vd2zzWf0=
@@ -10,8 +12,8 @@ github.com/alecthomas/colour v0.0.0-20160524082231-60882d9e2721 h1:JHZL0hZKJ1VEN
github.com/alecthomas/colour v0.0.0-20160524082231-60882d9e2721/go.mod h1:QO9JBoKquHd+jz9nshCh40fOfO+JzsoXy8qTHF68zU0= github.com/alecthomas/colour v0.0.0-20160524082231-60882d9e2721/go.mod h1:QO9JBoKquHd+jz9nshCh40fOfO+JzsoXy8qTHF68zU0=
github.com/alecthomas/kong v0.2.1-0.20190708041108-0548c6b1afae/go.mod h1:+inYUSluD+p4L8KdviBSgzcqEjUQOfC5fQDRFuc36lI= github.com/alecthomas/kong v0.2.1-0.20190708041108-0548c6b1afae/go.mod h1:+inYUSluD+p4L8KdviBSgzcqEjUQOfC5fQDRFuc36lI=
github.com/alecthomas/repr v0.0.0-20180818092828-117648cd9897/go.mod h1:xTS7Pm1pD1mvyM075QCDSRqH6qRLXylzS24ZTpRiSzQ= github.com/alecthomas/repr v0.0.0-20180818092828-117648cd9897/go.mod h1:xTS7Pm1pD1mvyM075QCDSRqH6qRLXylzS24ZTpRiSzQ=
github.com/alecthomas/repr v0.5.2 h1:SU73FTI9D1P5UNtvseffFSGmdNci/O6RsqzeXJtP0Qs= github.com/alecthomas/repr v0.5.3 h1:Ebk3yZ0kvrHC7TkTHLDJGDq1LxKeu9sQBQREFMcesS8=
github.com/alecthomas/repr v0.5.2/go.mod h1:Fr0507jx4eOXV7AlPV6AVZLYrLIuIeSOWtW57eE/O/4= github.com/alecthomas/repr v0.5.3/go.mod h1:Fr0507jx4eOXV7AlPV6AVZLYrLIuIeSOWtW57eE/O/4=
github.com/aymanbagabas/go-osc52/v2 v2.0.1 h1:HwpRHbFMcZLEVr42D4p7XBqjyuxQH5SMiErDT4WkJ2k= github.com/aymanbagabas/go-osc52/v2 v2.0.1 h1:HwpRHbFMcZLEVr42D4p7XBqjyuxQH5SMiErDT4WkJ2k=
github.com/aymanbagabas/go-osc52/v2 v2.0.1/go.mod h1:uYgXzlJ7ZpABp8OJ+exZzJJhRNQ2ASbcXHWsFqH8hp8= github.com/aymanbagabas/go-osc52/v2 v2.0.1/go.mod h1:uYgXzlJ7ZpABp8OJ+exZzJJhRNQ2ASbcXHWsFqH8hp8=
github.com/charmbracelet/bubbles v1.0.0 h1:12J8/ak/uCZEMQ6KU7pcfwceyjLlWsDLAxB5fXonfvc= github.com/charmbracelet/bubbles v1.0.0 h1:12J8/ak/uCZEMQ6KU7pcfwceyjLlWsDLAxB5fXonfvc=
@@ -151,8 +153,6 @@ github.com/rivo/uniseg v0.4.7 h1:WUdvkW8uEhrYfLC4ZzdpI2ztxP1I582+49Oc5Mq64VQ=
github.com/rivo/uniseg v0.4.7/go.mod h1:FN3SvrM+Zdj16jyLfmOkMNblXMcoc8DfTHruCPUcx88= github.com/rivo/uniseg v0.4.7/go.mod h1:FN3SvrM+Zdj16jyLfmOkMNblXMcoc8DfTHruCPUcx88=
github.com/rogpeppe/go-internal v1.13.1 h1:KvO1DLK/DRN07sQ1LQKScxyZJuNnedQ5/wKSR38lUII= github.com/rogpeppe/go-internal v1.13.1 h1:KvO1DLK/DRN07sQ1LQKScxyZJuNnedQ5/wKSR38lUII=
github.com/rogpeppe/go-internal v1.13.1/go.mod h1:uMEvuHeurkdAXX61udpOXGD/AzZDWNMNyH2VO9fmH0o= github.com/rogpeppe/go-internal v1.13.1/go.mod h1:uMEvuHeurkdAXX61udpOXGD/AzZDWNMNyH2VO9fmH0o=
github.com/seeruk/go-wordwrap v0.0.0-20191208221741-14ec4aac9550 h1:C3CfUXH/qmWuQFRqnPm3Sx8PFxa+pqACjhV5CaNO8pw=
github.com/seeruk/go-wordwrap v0.0.0-20191208221741-14ec4aac9550/go.mod h1:Sl541M2Em6rRG3V9WObycR7MYFZiERVkd/TJg0Gt0U4=
github.com/sergi/go-diff v1.0.0 h1:Kpca3qRNrduNnOQeazBd0ysaKrUJiIuISHxogkT9RPQ= github.com/sergi/go-diff v1.0.0 h1:Kpca3qRNrduNnOQeazBd0ysaKrUJiIuISHxogkT9RPQ=
github.com/sergi/go-diff v1.0.0/go.mod h1:0CfEIISq7TuYL3j771MWULgwwjU+GofnZX9QAmXWZgo= github.com/sergi/go-diff v1.0.0/go.mod h1:0CfEIISq7TuYL3j771MWULgwwjU+GofnZX9QAmXWZgo=
github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=

View File

@@ -194,18 +194,18 @@ func (cluster *Cluster) getDefaultOptions() []elasticsearch.Option {
} }
func (cluster *Cluster) getTransport() elastictransport.Option { func (cluster *Cluster) getTransport() elastictransport.Option {
transport := &http.Transport{ transport := new(http.Transport{
TLSClientConfig: &tls.Config{InsecureSkipVerify: true}, TLSClientConfig: &tls.Config{InsecureSkipVerify: true},
} })
if cluster.DebugHTTP { if cluster.DebugHTTP {
return elastictransport.WithTransport( return elastictransport.WithTransport(
&DebugTransport{Transport: transport}, new(DebugTransport{Transport: transport}),
) )
} }
return elastictransport.WithTransport( return elastictransport.WithTransport(
&CompatibilityTransport{Transport: transport}, new(CompatibilityTransport{Transport: transport}),
) )
} }

View File

@@ -22,13 +22,14 @@ import (
"os" "os"
"path/filepath" "path/filepath"
"reflect" "reflect"
"time"
"github.com/alecthomas/repr" "github.com/alecthomas/repr"
"gopkg.in/yaml.v3" "gopkg.in/yaml.v3"
) )
const ( const (
Version string = `v0.0.26` Version string = `v0.0.27`
) )
var ( var (
@@ -68,6 +69,13 @@ type Config struct {
Retention string // index template create: -r Retention string // index template create: -r
Rollover bool // index template create: -R Rollover bool // index template create: -R
Missing bool // index copy: -m
Sync bool // index copy: -S
RequestsPerSecond float64 // index copy: -R
Timeout time.Duration // index copy: -T
Refresh bool // index copy: -r
SourceIndices []string // index copy: -s
From, To, MaxItems int // search: flags From, To, MaxItems int // search: flags
Filter []string // search: -F Filter []string // search: -F
Path string // search+doc sh: -p Path string // search+doc sh: -p
@@ -90,6 +98,7 @@ type Config struct {
Force bool // ccr follower renew: -f Force bool // ccr follower renew: -f
DebugHTTP bool // root: --debug-http DebugHTTP bool // root: --debug-http
DebugGoRoutines bool // root: --debug-goroutines
Separator string // role diff: -s Separator string // role diff: -s
NotDeployed bool // role diff: -n NotDeployed bool // role diff: -n
Undefined bool // role diff: -u Undefined bool // role diff: -u
@@ -98,10 +107,12 @@ type Config struct {
// rollover // rollover
MaxAge string MaxAge string
MaxDocs, MaxShardSize, MaxShardDocs int // roll over MaxDocs int64 // roll over, plus others
MaxShardSize, MaxShardDocs int // roll over
DryRun bool // rollover: -n DryRun bool // rollover: -n
Tag string // api ls: -t Tag string // api ls: -t
HumanCat bool // api repl: -H
Ilm Ilm // ilm create Ilm Ilm // ilm create
@@ -112,7 +123,7 @@ type Config struct {
} }
func NewConfig() *Config { func NewConfig() *Config {
return &Config{Clusters: map[string]*Cluster{}} return new(Config{Clusters: map[string]*Cluster{}})
} }
func getDefaultPath() string { func getDefaultPath() string {
@@ -215,7 +226,7 @@ func (conf *Config) LoadConfig() error {
return fmt.Errorf("failed to read config file: %w", err) return fmt.Errorf("failed to read config file: %w", err)
} }
newconf := &Config{} newconf := new(Config{})
err = yaml.Unmarshal(data, newconf) err = yaml.Unmarshal(data, newconf)
if err != nil { if err != nil {

View File

@@ -28,6 +28,7 @@ import (
"io" "io"
"log" "log"
"log/slog" "log/slog"
"maps"
"net/http" "net/http"
"os" "os"
"os/exec" "os/exec"
@@ -45,22 +46,25 @@ import (
) )
const ( const (
intro = `Input format: verb path [data]" intro = `# Input format: verb path [data]"
#
Example: # Example:
#
post /yourindex/_ccr/pause_follow # post /yourindex/_ccr/pause_follow
put /yourindex/_settings {"number_of_replicas": 1} # put /yourindex/_settings {"number_of_replicas": 1}
#
You can also put multiline JSON after the path like: # You can also put multiline JSON after the path like:
#
put /yourindex/_settings # put /yourindex/_settings
{ # {
"number_of_replicas": 1 # "number_of_replicas": 1
} # }
#
If you do NOT supply a JSON in the first line, you need to hit ENTER # If you do NOT supply a JSON in the first line, you need to hit ENTER
twice to complete.` # twice to complete.
#
# Supply the flag --human-readable-cat, -H to view /_cat API calls in
# human readable form.`
) )
// holds an API operation via go-openapi/spec // holds an API operation via go-openapi/spec
@@ -143,15 +147,19 @@ func ApiRepl(conf *cfg.Config) error {
fmt.Printf("failed to call API: %s\n", esErrorString(err)) fmt.Printf("failed to call API: %s\n", esErrorString(err))
} }
if conf.HumanCat && strings.HasPrefix(parts[1], "/_cat") {
fmt.Println(string(raw))
} else {
pageJsonOutput(conf, raw) pageJsonOutput(conf, raw)
} }
}
//nolint:nilerr //nolint:nilerr
return nil return nil
} }
func pageJsonOutput(conf *cfg.Config, raw []byte) { func pageJsonOutput(conf *cfg.Config, raw []byte) {
tmpconf := &cfg.Config{HaveJQ: conf.HaveJQ} tmpconf := new(cfg.Config{HaveJQ: conf.HaveJQ})
if conf.Pager != "" { if conf.Pager != "" {
tmpconf.HaveJQ = false tmpconf.HaveJQ = false
@@ -199,16 +207,16 @@ func CallAPI(conf *cfg.Config, verb, path, data string) ([]byte, error) {
verb = strings.ToUpper(verb) verb = strings.ToUpper(verb)
// we're using port-forwards anyway // we're using port-forwards anyway
noVerifyTransport := &http.Transport{ noVerifyTransport := new(http.Transport{
TLSClientConfig: &tls.Config{InsecureSkipVerify: true}, TLSClientConfig: new(tls.Config{InsecureSkipVerify: true}),
} })
client := &http.Client{Transport: noVerifyTransport} client := new(http.Client{Transport: noVerifyTransport})
if conf.DebugHTTP { if conf.DebugHTTP {
client = &http.Client{ client = new(http.Client{
Transport: &cfg.DebugTransport{ Transport: new(cfg.DebugTransport{
Transport: noVerifyTransport}} Transport: noVerifyTransport})})
} }
req, err := http.NewRequest(verb, conf.DefaultCluster.Uri+path, bytes.NewBuffer([]byte(data))) req, err := http.NewRequest(verb, conf.DefaultCluster.Uri+path, bytes.NewBuffer([]byte(data)))
@@ -216,8 +224,10 @@ func CallAPI(conf *cfg.Config, verb, path, data string) ([]byte, error) {
return nil, err return nil, err
} }
if !conf.HumanCat || (conf.HumanCat && !strings.HasPrefix(path, "/_cat")) {
req.Header.Add("Content-Type", "application/json") req.Header.Add("Content-Type", "application/json")
req.Header.Add("Accept", "application/json") req.Header.Add("Accept", "application/json")
}
// make sure we have got all we need // make sure we have got all we need
if err := conf.DefaultCluster.CheckAuth(); err != nil { if err := conf.DefaultCluster.CheckAuth(); err != nil {
@@ -271,7 +281,8 @@ func prettyfiJson(conf *cfg.Config, raw []byte) (string, error) {
err := json.Indent(&pretty, raw, "", "\t") err := json.Indent(&pretty, raw, "", "\t")
if err != nil { if err != nil {
return "", fmt.Errorf("json parse error: %w", err) //nolint:nilerr
return string(raw), nil
} }
return pretty.String(), nil return pretty.String(), nil
@@ -317,7 +328,7 @@ func ApiList(conf *cfg.Config, pattern string) error {
filter := regexp.MustCompile(pattern) filter := regexp.MustCompile(pattern)
table := printer.NewTable(conf, 4, 0) table := printer.NewTable(conf).WithSize(4, 0)
table.Addheaders("path", "http verb", "tag", "description") table.Addheaders("path", "http verb", "tag", "description")
for path, item := range assets.OpenAPI.Spec().Paths.Paths { for path, item := range assets.OpenAPI.Spec().Paths.Paths {
@@ -356,16 +367,7 @@ func ApiList(conf *cfg.Config, pattern string) error {
func ApiPathNames() []string { func ApiPathNames() []string {
assets.LoadAssetOpenApi() assets.LoadAssetOpenApi()
paths := make([]string, len(assets.OpenAPI.Spec().Paths.Paths)) return slices.Collect(maps.Keys(assets.OpenAPI.Spec().Paths.Paths))
idx := 0
for path := range assets.OpenAPI.Spec().Paths.Paths {
paths[idx] = path
idx++
}
return paths
} }
func ApiShow(conf *cfg.Config, showpath, verb string) error { func ApiShow(conf *cfg.Config, showpath, verb string) error {
@@ -526,7 +528,7 @@ func getApiExample(op *Op) string {
// otherwise showpath+verb have to match precisely. // otherwise showpath+verb have to match precisely.
func matchOperation(showpath, verb string) (*Op, error) { func matchOperation(showpath, verb string) (*Op, error) {
ops := []*Op{} ops := []*Op{}
op := &Op{} op := new(Op{})
var found bool var found bool

View File

@@ -57,7 +57,7 @@ func CcrStatus(conf *cfg.Config, leader, follower string) error {
res, err := conf.Clusters[alias].ES().Cat.Indices(). res, err := conf.Clusters[alias].ES().Cat.Indices().
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
return fmt.Errorf("failed to get indicies on %s: %w", alias, esErrorString(err)) return fmt.Errorf("failed to get indices on %s: %w", alias, esErrorString(err))
} }
indices[alias] = make(map[string]*types.IndicesRecord, len(res)) indices[alias] = make(map[string]*types.IndicesRecord, len(res))
@@ -111,7 +111,7 @@ func CcrRemoteInfo(conf *cfg.Config, index string) error {
mode = "leader" mode = "leader"
} }
table := printer.NewTable(conf, 2, 5) table := printer.NewTable(conf).WithSize(2, 5)
table.Addheaders("ccr remote property", "value") table.Addheaders("ccr remote property", "value")
table.Entries = [][]any{ table.Entries = [][]any{

View File

@@ -164,7 +164,7 @@ func CcrFollowerShow(conf *cfg.Config, index string) error {
follower := res.Indices[0].Shards[0] follower := res.Indices[0].Shards[0]
table := printer.NewTable(conf, 2, 9) table := printer.NewTable(conf).WithSize(2, 9)
table.Addheaders("ccr follower property", "value") table.Addheaders("ccr follower property", "value")
table.Entries = [][]any{ table.Entries = [][]any{

View File

@@ -17,6 +17,10 @@ along with this program. If not, see <http://www.gnu.org/licenses/>.
package es package es
import ( import (
// "encoding/json/jsontext"
// "encoding/json/v2"
"encoding/json" "encoding/json"
"fmt" "fmt"
@@ -49,11 +53,22 @@ func getHealthReport(conf *cfg.Config) (*HealthReport, error) {
return nil, err return nil, err
} }
report := HealthReport{} report := new(HealthReport{})
// FIXME: use this once jsonv2 is no more experimental it already
// builds and works like intended, but golangci-lint doesn't
// recognize it with: go: unknown GOEXPERIMENT jsonv2
//
// if err := json.UnmarshalDecode(
// jsontext.NewDecoder(
// bytes.NewBuffer(raw)),
// &report); err != nil {
// return nil, fmt.Errorf("failed to unmarshal healthreport response: %w", err)
// }
if err := json.Unmarshal(raw, &report); err != nil { if err := json.Unmarshal(raw, &report); err != nil {
return nil, fmt.Errorf("failed to unmarshal healthreport response: %w", err) return nil, fmt.Errorf("failed to unmarshal healthreport response: %w", err)
} }
return &report, nil return report, nil
} }

View File

@@ -58,7 +58,7 @@ func ClusterList(conf *cfg.Config) error {
wg.Wait() wg.Wait()
table := printer.NewTable(conf, 5, len(conf.Clusters)) table := printer.NewTable(conf).WithSize(5, len(conf.Clusters))
table.Addheaders("cluster", "uri", "reachable", "current", "error") table.Addheaders("cluster", "uri", "reachable", "current", "error")
idx := 0 idx := 0
@@ -101,23 +101,39 @@ func getClusterStatus(conf *cfg.Config) (*apiResponse, error) {
es := conf.DefaultCluster.ES() es := conf.DefaultCluster.ES()
responses := make(chan apiResponse, gocount) responses := make(chan apiResponse, gocount)
wg := &sync.WaitGroup{} wg := new(sync.WaitGroup{})
wg.Add(gocount) wg.Go(func() {
go getApiData(conf, es, wg, responses, "health") getApiData(conf, es, responses, "health")
go getApiData(conf, es, wg, responses, "healthreport") })
go getApiData(conf, es, wg, responses, "info")
go getApiData(conf, es, wg, responses, "ccr") wg.Go(func() {
go getApiData(conf, es, wg, responses, "indices") getApiData(conf, es, responses, "healthreport")
go getApiData(conf, es, wg, responses, "tasks") })
wg.Go(func() {
getApiData(conf, es, responses, "info")
})
wg.Go(func() {
getApiData(conf, es, responses, "ccr")
})
wg.Go(func() {
getApiData(conf, es, responses, "indices")
})
wg.Go(func() {
getApiData(conf, es, responses, "tasks")
})
if conf.Verbose { if conf.Verbose {
go getApiData(conf, es, wg, responses, "stats") getApiData(conf, es, responses, "stats")
} }
wg.Wait() wg.Wait()
all := apiResponse{} all := new(apiResponse{})
var err error var err error
@@ -144,7 +160,7 @@ func getClusterStatus(conf *cfg.Config) (*apiResponse, error) {
} }
} }
return &all, err return all, err
} }
func ClusterStatus(conf *cfg.Config) error { func ClusterStatus(conf *cfg.Config) error {
@@ -196,7 +212,7 @@ func ClusterStatus(conf *cfg.Config) error {
} }
} }
table := printer.NewTable(conf, 2, 7) table := printer.NewTable(conf).WithSize(2, 7)
table.Addheaders(conf.DefaultCluster.Name, "status") table.Addheaders(conf.DefaultCluster.Name, "status")
table.Entries = [][]any{ table.Entries = [][]any{
@@ -272,10 +288,10 @@ func gatherClusterStats(clusterstats *clusterstats.Response, table *printer.Tabl
} }
table.Entries = append(table.Entries, [][]any{ table.Entries = append(table.Entries, [][]any{
{"Indicies", clusterstats.Indices.Count}, {"indices", clusterstats.Indices.Count},
{"Docs", clusterstats.Indices.Docs.Count}, {"Docs", clusterstats.Indices.Docs.Count},
{"Total Size", printer.Bytes(clusterstats.Indices.Docs.TotalSizeInBytes)}, {"Total Size", printer.Bytes(clusterstats.Indices.Docs.TotalSizeInBytes)},
{"Total Queries", "%d", querycount}, {"Total Queries", querycount},
{"Shards Primaries", clusterstats.Indices.Shards.Primaries}, {"Shards Primaries", clusterstats.Indices.Shards.Primaries},
{"Shards Total", clusterstats.Indices.Shards.Total}, {"Shards Total", clusterstats.Indices.Shards.Total},
{"Storage", fmt.Sprintf( {"Storage", fmt.Sprintf(

View File

@@ -29,12 +29,12 @@ func ClusterRerouteMove(conf *cfg.Config, index string) error {
move := conf.DefaultCluster.ES().Cluster.Reroute() move := conf.DefaultCluster.ES().Cluster.Reroute()
commands := esdsl.NewCommand() commands := esdsl.NewCommand()
moveCommand := &types.CommandMoveAction{ moveCommand := new(types.CommandMoveAction{
Shard: conf.Shards, Shard: conf.Shards,
FromNode: conf.FromNode, FromNode: conf.FromNode,
ToNode: conf.ToNode, ToNode: conf.ToNode,
Index: index, Index: index,
} })
commands.CommandCaster().Move = moveCommand commands.CommandCaster().Move = moveCommand
@@ -52,11 +52,11 @@ func ClusterRerouteAllocateReplica(conf *cfg.Config, index string) error {
move := conf.DefaultCluster.ES().Cluster.Reroute() move := conf.DefaultCluster.ES().Cluster.Reroute()
commands := esdsl.NewCommand() commands := esdsl.NewCommand()
allocCommand := &types.CommandAllocateReplicaAction{ allocCommand := new(types.CommandAllocateReplicaAction{
Shard: conf.Shards, Shard: conf.Shards,
Node: conf.ToNode, Node: conf.ToNode,
Index: index, Index: index,
} })
commands.CommandCaster().AllocateReplica = allocCommand commands.CommandCaster().AllocateReplica = allocCommand
@@ -74,12 +74,12 @@ func ClusterRerouteCancel(conf *cfg.Config, index string) error {
move := conf.DefaultCluster.ES().Cluster.Reroute() move := conf.DefaultCluster.ES().Cluster.Reroute()
commands := esdsl.NewCommand() commands := esdsl.NewCommand()
cancelCommand := &types.CommandCancelAction{ cancelCommand := new(types.CommandCancelAction{
Shard: conf.Shards, Shard: conf.Shards,
Node: conf.ToNode, Node: conf.ToNode,
Index: index, Index: index,
AllowPrimary: &conf.AllowPrimary, AllowPrimary: &conf.AllowPrimary,
} })
commands.CommandCaster().Cancel = cancelCommand commands.CommandCaster().Cancel = cancelCommand
@@ -97,12 +97,12 @@ func ClusterRerouteAllocatePrimary(conf *cfg.Config, index string, stale bool) e
move := conf.DefaultCluster.ES().Cluster.Reroute() move := conf.DefaultCluster.ES().Cluster.Reroute()
commands := esdsl.NewCommand() commands := esdsl.NewCommand()
allocCommand := &types.CommandAllocatePrimaryAction{ allocCommand := new(types.CommandAllocatePrimaryAction{
Shard: conf.Shards, Shard: conf.Shards,
Node: conf.ToNode, Node: conf.ToNode,
Index: index, Index: index,
AcceptDataLoss: conf.AcceptDataLoss, AcceptDataLoss: conf.AcceptDataLoss,
} })
if stale { if stale {
commands.CommandCaster().AllocateStalePrimary = allocCommand commands.CommandCaster().AllocateStalePrimary = allocCommand

View File

@@ -38,7 +38,7 @@ func ClusterSettingsList(conf *cfg.Config) error {
return fmt.Errorf("failed to get cluster settings: %w", esErrorString(err)) return fmt.Errorf("failed to get cluster settings: %w", esErrorString(err))
} }
table := printer.NewTable(conf, 2, 0) table := printer.NewTable(conf).WithSize(2, 0)
table.Addheaders("setting", "value") table.Addheaders("setting", "value")
entries := [][]any{} entries := [][]any{}

View File

@@ -71,7 +71,7 @@ func checkClusterStatus(conf *cfg.Config, leader, follower string) bool {
status[cluster] = clusterHealth status[cluster] = clusterHealth
} }
table := printer.NewTable(conf, 3, 6) table := printer.NewTable(conf).WithSize(3, 6)
table.Addheaders("setting", "leader:"+leader, "follower:"+follower) table.Addheaders("setting", "leader:"+leader, "follower:"+follower)
@@ -89,7 +89,7 @@ func checkClusterStatus(conf *cfg.Config, leader, follower string) bool {
status[leader].ActivePrimaryShards, status[leader].ActivePrimaryShards,
status[follower].ActivePrimaryShards, status[follower].ActivePrimaryShards,
}, },
{"Indicies", {"indices",
len(status[leader].Indices), len(status[leader].Indices),
len(status[follower].Indices), len(status[follower].Indices),
}, },
@@ -151,7 +151,7 @@ func findIlmErrors(conf *cfg.Config, leader, follower string) bool {
if len(failed[cluster]) > 0 { if len(failed[cluster]) > 0 {
idx := 0 idx := 0
table := printer.NewTable(conf, 2, len(failed[cluster])) table := printer.NewTable(conf).WithSize(2, len(failed[cluster]))
table.Addheaders("ilm errors on "+which, "errors") table.Addheaders("ilm errors on "+which, "errors")
for name, count := range failed[cluster] { for name, count := range failed[cluster] {
@@ -172,7 +172,7 @@ func findIlmErrors(conf *cfg.Config, leader, follower string) bool {
return false return false
} }
// find unsynchronized indicies only present on leader // find unsynchronized indices only present on leader
func findIndicesOnlyOnLeader(conf *cfg.Config, indices ClusterIndices, leader, follower string) bool { func findIndicesOnlyOnLeader(conf *cfg.Config, indices ClusterIndices, leader, follower string) bool {
exclude := regexp.MustCompile(DefaultExclude) exclude := regexp.MustCompile(DefaultExclude)
if conf.Exclude != "" { if conf.Exclude != "" {
@@ -225,7 +225,7 @@ func findIndicesOnlyOnLeader(conf *cfg.Config, indices ClusterIndices, leader, f
} }
idx := 0 idx := 0
table := printer.NewTable(conf, 3, len(indexOnlyOnLeader)) table := printer.NewTable(conf).WithSize(3, len(indexOnlyOnLeader))
table.Addheaders("index only on leader", "size", "docscount") table.Addheaders("index only on leader", "size", "docscount")
for name, index := range indexOnlyOnLeader { for name, index := range indexOnlyOnLeader {
@@ -275,7 +275,7 @@ func findOrphanedIndices(conf *cfg.Config, indices ClusterIndices, leader, follo
} }
idx := 0 idx := 0
table := printer.NewTable(conf, 3, len(orphaned)) table := printer.NewTable(conf).WithSize(3, len(orphaned))
table.Addheaders("orphaned index on follower", "size", "docscount") table.Addheaders("orphaned index on follower", "size", "docscount")
for name, index := range orphaned { for name, index := range orphaned {
@@ -311,7 +311,7 @@ func findFailedFollowerIndices(conf *cfg.Config, indices ClusterIndices, followe
} }
idx := 0 idx := 0
table := printer.NewTable(conf, 3, len(red)) table := printer.NewTable(conf).WithSize(3, len(red))
table.Addheaders("red index on follower", "size", "docscount") table.Addheaders("red index on follower", "size", "docscount")
for name, index := range red { for name, index := range red {

View File

@@ -62,7 +62,7 @@ func DatastreamList(conf *cfg.Config) error {
} }
} }
table := printer.NewTable(conf, 7, size) table := printer.NewTable(conf).WithSize(7, size)
table.Addheaders("name", "ilm policy", "hidden", "system", "replicated", "generation", "timestamp field") table.Addheaders("name", "ilm policy", "hidden", "system", "replicated", "generation", "timestamp field")
for idx, datastream := range list { for idx, datastream := range list {
@@ -146,7 +146,7 @@ func DatastreamShow(conf *cfg.Config, dsname string) error {
return fmt.Errorf("failed to get data stream stats: %w", esErrorString(err)) return fmt.Errorf("failed to get data stream stats: %w", esErrorString(err))
} }
table := printer.NewTable(conf, 11, 2) table := printer.NewTable(conf).WithSize(11, 2)
table.Addheaders("data stream property", "value") table.Addheaders("data stream property", "value")
datastream := res.DataStreams[0] datastream := res.DataStreams[0]
@@ -176,7 +176,7 @@ func DatastreamShow(conf *cfg.Config, dsname string) error {
return err return err
} }
table = printer.NewTable(conf, 5, len(datastream.Indices)) table = printer.NewTable(conf).WithSize(5, len(datastream.Indices))
table.Addheaders("backing index name", "uuid", "prefer ilm", "ilm policy", "managed by") table.Addheaders("backing index name", "uuid", "prefer ilm", "ilm policy", "managed by")
for idx, index := range datastream.Indices { for idx, index := range datastream.Indices {
@@ -233,7 +233,7 @@ func DatastreamRollover(conf *cfg.Config, ds string) error {
return fmt.Errorf("failed to rollover data stream: %w", esErrorString(err)) return fmt.Errorf("failed to rollover data stream: %w", esErrorString(err))
} }
table := printer.NewTable(conf, 2, 5) table := printer.NewTable(conf).WithSize(2, 5)
table.Addheaders("rollover response", "value") table.Addheaders("rollover response", "value")
table.Entries = [][]any{ table.Entries = [][]any{
{"acknowledged", res.Acknowledged}, {"acknowledged", res.Acknowledged},

View File

@@ -101,7 +101,7 @@ func DocDelete(conf *cfg.Config, queries []string) error {
return nil return nil
} }
req := &deletebyquery.Request{} req := new(deletebyquery.Request{})
if len(queries) == 0 && conf.All { if len(queries) == 0 && conf.All {
req.Query = esdsl.NewMatchAllQuery().QueryCaster() req.Query = esdsl.NewMatchAllQuery().QueryCaster()

View File

@@ -22,6 +22,8 @@ import (
"errors" "errors"
"fmt" "fmt"
"log/slog" "log/slog"
"maps"
"slices"
"strings" "strings"
"codeberg.org/scip/esctl/pkg/cfg" "codeberg.org/scip/esctl/pkg/cfg"
@@ -62,13 +64,7 @@ func IlmNames(conf *cfg.Config) ([]string, error) {
return nil, fmt.Errorf("failed to get ilm policies: %w", esErrorString(err)) return nil, fmt.Errorf("failed to get ilm policies: %w", esErrorString(err))
} }
names := make([]string, len(res)) names := slices.Collect(maps.Keys(res))
idx := 0
for name := range res {
names[idx] = name
idx++
}
return names, nil return names, nil
} }
@@ -89,7 +85,7 @@ func IlmList(conf *cfg.Config, pattern string) error {
repr.Println(res) repr.Println(res)
} }
table := printer.NewTable(conf, 5, 0) table := printer.NewTable(conf).WithSize(5, 0)
table.Addheaders("ilm policy", "hot", "warm", "frozen", "delete") table.Addheaders("ilm policy", "hot", "warm", "frozen", "delete")
for name, ilm := range res { for name, ilm := range res {
@@ -127,7 +123,7 @@ func IlmShow(conf *cfg.Config, policy string) error {
return IlmShowTree(conf, ilm.Policy) return IlmShowTree(conf, ilm.Policy)
} }
table := printer.NewTable(conf, 2, 5) table := printer.NewTable(conf).WithSize(2, 5)
table.Addheaders("ilm policy setting", "value") table.Addheaders("ilm policy setting", "value")
table.Entries = [][]any{ table.Entries = [][]any{
@@ -146,7 +142,7 @@ func IlmShow(conf *cfg.Config, policy string) error {
func IlmShowTree(conf *cfg.Config, ilm types.IlmPolicy) error { func IlmShowTree(conf *cfg.Config, ilm types.IlmPolicy) error {
indent := "" indent := ""
table := printer.NewTable(conf, 4, 0) table := printer.NewTable(conf).WithSize(4, 0)
table.Addheaders("phase", "min age", "min size", "snapshot repo") table.Addheaders("phase", "min age", "min size", "snapshot repo")
for _, phase := range IlmPhaseOrder { for _, phase := range IlmPhaseOrder {
@@ -284,15 +280,17 @@ func IlmExplain(conf *cfg.Config, index string) error {
ilm := explain.(*types.LifecycleExplainManaged) ilm := explain.(*types.LifecycleExplainManaged)
table := printer.NewTable(conf, 2, 0) table := printer.NewTable(conf).WithSize(2, 0)
table.Addheaders("ilm status field", "value") table.Addheaders("ilm status field", "value")
info := "" info := ""
if len(ilm.StepInfo["reason"]) > 0 {
err = json.Unmarshal(ilm.StepInfo["reason"], &info) err = json.Unmarshal(ilm.StepInfo["reason"], &info)
if err != nil { if err != nil {
return fmt.Errorf("failed to unmarshal step info: %w", err) return fmt.Errorf("failed to unmarshal step info: %w", err)
} }
}
table.Entries = [][]any{ table.Entries = [][]any{
{"index", ilm.Index}, {"index", ilm.Index},
@@ -303,8 +301,11 @@ func IlmExplain(conf *cfg.Config, index string) error {
{"phase", *ilm.Phase}, {"phase", *ilm.Phase},
{"phase execution", ilmPhaseString(ilm.PhaseExecution.PhaseDefinition, false)}, {"phase execution", ilmPhaseString(ilm.PhaseExecution.PhaseDefinition, false)},
{"step", *ilm.Step}, {"step", *ilm.Step},
{"failed step", *ilm.FailedStep}, }
{"failed step retry count", ilm.FailedStepRetryCount},
if ilm.FailedStep != nil {
table.AddRow("failed step", *ilm.FailedStep)
table.AddRow("failed step retry count", ilm.FailedStepRetryCount)
} }
if err := table.Print(); err != nil { if err := table.Print(); err != nil {
@@ -351,7 +352,7 @@ func IlmCreate(conf *cfg.Config, policyname string) error {
var actions types.IlmActionsVariant = esdsl.NewIlmActions() var actions types.IlmActionsVariant = esdsl.NewIlmActions()
rollover := &types.RolloverAction{} rollover := new(types.RolloverAction{})
haveroll := false haveroll := false
if policy != nil { if policy != nil {
@@ -513,8 +514,8 @@ func IlmCreate(conf *cfg.Config, policyname string) error {
phases.PhasesCaster().Delete = policy.Phases.Delete phases.PhasesCaster().Delete = policy.Phases.Delete
} }
put := &putlifecycle.Request{} put := new(putlifecycle.Request{})
newpolicy := &types.IlmPolicy{} newpolicy := new(types.IlmPolicy{})
newpolicy.IlmPolicyCaster().Phases = *phases.PhasesCaster() newpolicy.IlmPolicyCaster().Phases = *phases.PhasesCaster()
put.Policy = newpolicy put.Policy = newpolicy

View File

@@ -105,7 +105,7 @@ func IlmForecastList(conf *cfg.Config, filter string) error {
headers = append(headers, "ilm policy") headers = append(headers, "ilm policy")
} }
table := printer.NewTable(conf, len(headers), 0) table := printer.NewTable(conf).WithSize(len(headers), 0)
table.Addheaders(headers...) table.Addheaders(headers...)
for _, phase := range phaseData { for _, phase := range phaseData {
@@ -185,12 +185,19 @@ func virtualAge(phase *PhaseData) time.Duration {
// Retrieve all index, ilm-explain and ilm-policies in parallel // Retrieve all index, ilm-explain and ilm-policies in parallel
func getIlmPhaseData(conf *cfg.Config) ([]PhaseData, error) { func getIlmPhaseData(conf *cfg.Config) ([]PhaseData, error) {
responses := make(chan apiResponse, 3) responses := make(chan apiResponse, 3)
wg := &sync.WaitGroup{} wg := new(sync.WaitGroup{})
wg.Add(3)
go getApiData(conf, conf.DefaultCluster.ES(), wg, responses, "indicesbytes") wg.Go(func() {
go getApiData(conf, conf.DefaultCluster.ES(), wg, responses, "explain") getApiData(conf, conf.DefaultCluster.ES(), responses, "indicesbytes")
go getApiData(conf, conf.DefaultCluster.ES(), wg, responses, "policies") })
wg.Go(func() {
getApiData(conf, conf.DefaultCluster.ES(), responses, "explain")
})
wg.Go(func() {
getApiData(conf, conf.DefaultCluster.ES(), responses, "policies")
})
wg.Wait() wg.Wait()
@@ -331,7 +338,7 @@ func findNextPhase(policy types.IlmPolicy, currentPhase string) *NextPhase {
// phase list to determine which comes next // phase list to determine which comes next
phases, start := registerPhases(policy, currentPhase) phases, start := registerPhases(policy, currentPhase)
nextPhase := &NextPhase{} nextPhase := new(NextPhase{})
// finally determine which phase comes next // finally determine which phase comes next
// exception: hot, where we look for rollover rules // exception: hot, where we look for rollover rules

View File

@@ -18,6 +18,8 @@ package es
import ( import (
"context" "context"
"encoding/json"
"errors"
"fmt" "fmt"
"log/slog" "log/slog"
"regexp" "regexp"
@@ -28,8 +30,11 @@ import (
"codeberg.org/scip/esctl/pkg/cfg" "codeberg.org/scip/esctl/pkg/cfg"
"codeberg.org/scip/esctl/pkg/printer" "codeberg.org/scip/esctl/pkg/printer"
"codeberg.org/scip/mapmap"
"github.com/charmbracelet/lipgloss"
"github.com/elastic/go-elasticsearch/v9/typedapi/cat/indices" "github.com/elastic/go-elasticsearch/v9/typedapi/cat/indices"
"github.com/elastic/go-elasticsearch/v9/typedapi/esdsl" "github.com/elastic/go-elasticsearch/v9/typedapi/esdsl"
"github.com/elastic/go-elasticsearch/v9/typedapi/types"
"github.com/elastic/go-elasticsearch/v9/typedapi/types/enums/healthstatus" "github.com/elastic/go-elasticsearch/v9/typedapi/types/enums/healthstatus"
) )
@@ -38,7 +43,7 @@ func IndexNames(conf *cfg.Config) ([]string, error) {
res, err := conf.DefaultCluster.ES().Cat.Indices(). res, err := conf.DefaultCluster.ES().Cat.Indices().
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
return nil, fmt.Errorf("failed to get indicies: %w", esErrorString(err)) return nil, fmt.Errorf("failed to get indices: %w", esErrorString(err))
} }
indices := make([]string, len(res)) indices := make([]string, len(res))
@@ -50,41 +55,29 @@ func IndexNames(conf *cfg.Config) ([]string, error) {
} }
func filterIndices(conf *cfg.Config, list indices.Response) indices.Response { func filterIndices(conf *cfg.Config, list indices.Response) indices.Response {
// apply partials filter first var filter *regexp.Regexp
selectedlist := indices.Response{}
for _, index := range list { if len(conf.Filter) > 0 {
filter = regexp.MustCompile(conf.Filter[0])
}
return mapmap.NewSlicer(list).MapSliceValuesImmutable(func(index types.IndicesRecord) bool {
if !conf.Partials && strings.HasPrefix(*index.Index, "partial-") { if !conf.Partials && strings.HasPrefix(*index.Index, "partial-") {
continue return false
} }
if !conf.Hidden && strings.HasPrefix(*index.Index, ".") { if !conf.Hidden && strings.HasPrefix(*index.Index, ".") {
continue return false
} }
if strings.HasPrefix(*index.Index, ".ds-") { if len(conf.Filter) > 0 {
// ignore data stream backing indicies if !filter.MatchString(*index.Index) {
continue return false
}
selectedlist = append(selectedlist, index)
}
if len(conf.Filter) == 0 {
return selectedlist
}
// we support just one filter here, for now
filter := *regexp.MustCompile(conf.Filter[0])
newlist := indices.Response{}
for _, index := range selectedlist {
if filter.MatchString(*index.Index) {
newlist = append(newlist, index)
} }
} }
return newlist return true
})
} }
func IndexList(conf *cfg.Config) error { func IndexList(conf *cfg.Config) error {
@@ -96,10 +89,10 @@ func IndexList(conf *cfg.Config) error {
res, err := cat.Do(context.Background()) res, err := cat.Do(context.Background())
if err != nil { if err != nil {
return fmt.Errorf("failed to get indicies: %w", esErrorString(err)) return fmt.Errorf("failed to get indices: %w", esErrorString(err))
} }
slog.Debug("ES result", "indicies", res) slog.Debug("ES result", "indices", res)
list := filterIndices(conf, res) list := filterIndices(conf, res)
@@ -111,7 +104,7 @@ func IndexList(conf *cfg.Config) error {
} }
} }
table := printer.NewTable(conf, 3, size) table := printer.NewTable(conf).WithSize(3, size)
table.Addheaders("name", "size", "docscount") table.Addheaders("name", "size", "docscount")
for idx, index := range list { for idx, index := range list {
@@ -147,7 +140,7 @@ func IndexShow(conf *cfg.Config, indexpattern string) error {
idx++ idx++
} }
table := printer.NewTable(conf, 2, 7) table := printer.NewTable(conf).WithSize(2, 7)
table.Addheaders("index property", "value") table.Addheaders("index property", "value")
ts, err := strconv.ParseInt(index.Settings.Index.CreationDate.(string), 10, 64) ts, err := strconv.ParseInt(index.Settings.Index.CreationDate.(string), 10, 64)
@@ -270,7 +263,7 @@ func IndexFields(conf *cfg.Config, index string) error {
return fmt.Errorf("failed to retrieve field capabilties: %w", esErrorString(err)) return fmt.Errorf("failed to retrieve field capabilties: %w", esErrorString(err))
} }
table := printer.NewTable(conf, 5, 0) table := printer.NewTable(conf).WithSize(5, 0)
table.Addheaders("field", "type", "searchable", "aggretable", "metadata") table.Addheaders("field", "type", "searchable", "aggretable", "metadata")
idx := 0 idx := 0
@@ -306,3 +299,79 @@ func IndexFields(conf *cfg.Config, index string) error {
return nil return nil
} }
type Diskusage struct {
Total int64 `json:"total_in_bytes"`
Points int64 `json:"points_in_bytes"`
Norms int64 `json:"norms_in_bytes"`
TermVectors int64 `json:"term_vectors_in_bytes"`
KnnVectors int64 `json:"knn_vectors_in_bytes"`
BloomFilter int64 `json:"bloom_filter_in_bytes"`
}
type IndexDiskUsage struct {
AllFields Diskusage `json:"all_fields"`
Fields map[string]Diskusage `json:"fields"`
}
type ResIndexDiskUsage map[string]IndexDiskUsage
func IndexDiskusage(conf *cfg.Config, index string) error {
var bold = lipgloss.NewStyle().Bold(true)
res, err := conf.DefaultCluster.ES().Indices.DiskUsage(index).
RunExpensiveTasks(true).
Do(context.Background())
if err != nil {
return fmt.Errorf("failed to retrieve index disk usage: %w", esErrorString(err))
}
duRes := ResIndexDiskUsage{}
if err := json.Unmarshal(res, &duRes); err != nil {
return fmt.Errorf("failed to unmarshal disk usage response: %w", err)
}
diskusage, exists := duRes[index]
if !exists {
return errors.New("no disk usage reported for index")
}
table := printer.NewTable(conf).
WithHeaders("field", "bloom filter", "norms", "points", "term vectors", "knn vectors", "total")
for name, field := range diskusage.Fields {
if strings.HasPrefix(name, "_") || strings.HasSuffix(name, ".keyword") {
continue
}
table.AddRow(
name,
printer.Bytes(field.BloomFilter),
printer.Bytes(field.Norms),
printer.Bytes(field.Points),
printer.Bytes(field.TermVectors),
printer.Bytes(field.KnnVectors),
printer.Bytes(field.Total),
)
}
all := diskusage.AllFields
table.Sort()
table.AddRowLate(
bold.Render("Summary"),
printer.Bytes(all.BloomFilter),
printer.Bytes(all.Norms),
printer.Bytes(all.Points),
printer.Bytes(all.TermVectors),
printer.Bytes(all.KnnVectors),
printer.Bytes(all.Total),
)
if err := table.Print(); err != nil {
return err
}
return nil
}

View File

@@ -72,7 +72,7 @@ func IndexAliasList(conf *cfg.Config) error {
} }
} }
table := printer.NewTable(conf, 2, len(aliaslist)) table := printer.NewTable(conf).WithSize(2, len(aliaslist))
table.Addheaders("index", "alias") table.Addheaders("index", "alias")
idx := 0 idx := 0
@@ -113,7 +113,7 @@ func IndexAliasRollover(conf *cfg.Config, alias string) error {
return fmt.Errorf("failed to rollover index alias: %w", esErrorString(err)) return fmt.Errorf("failed to rollover index alias: %w", esErrorString(err))
} }
table := printer.NewTable(conf, 2, 5) table := printer.NewTable(conf).WithSize(2, 5)
table.Addheaders("rollover response", "value") table.Addheaders("rollover response", "value")
table.Entries = [][]any{ table.Entries = [][]any{
{"acknowledged", res.Acknowledged}, {"acknowledged", res.Acknowledged},

83
pkg/es/index_copy.go Normal file
View File

@@ -0,0 +1,83 @@
/*
Copyright © 2026 Thomas von Dein
This program is free software: you can redistribute it and/or modify
it under the terms of the GNU General Public License as published by
the Free Software Foundation, either version 3 of the License, or
(at your option) any later version.
This program is distributed in the hope that it will be useful,
but WITHOUT ANY WARRANTY; without even the implied warranty of
MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
GNU General Public License for more details.
You should have received a copy of the GNU General Public License
along with this program. If not, see <http://www.gnu.org/licenses/>.
*/
package es
import (
"context"
"fmt"
"time"
"codeberg.org/scip/esctl/pkg/cfg"
"codeberg.org/scip/esctl/pkg/printer"
"github.com/elastic/go-elasticsearch/v9/typedapi/esdsl"
"github.com/elastic/go-elasticsearch/v9/typedapi/types/enums/conflicts"
)
func IndexCopy(conf *cfg.Config) error {
copy := conf.DefaultCluster.ES().Reindex()
if conf.Force {
copy.Conflicts(conflicts.Proceed)
}
if conf.RequestsPerSecond != 0 {
copy.RequestsPerSecond(fmt.Sprintf("%.2f", conf.RequestsPerSecond))
}
if conf.MaxDocs > 0 {
copy.MaxDocs(conf.MaxDocs)
}
if conf.Timeout > 0 {
copy.Timeout(formatDuration(conf.Timeout))
}
if conf.Wait {
copy.WaitForActiveShards("all")
}
if conf.Refresh {
copy.Refresh(true)
}
copy.Source(esdsl.NewReindexSource().Index(conf.SourceIndices...))
copy.Dest(esdsl.NewReindexDestination().Index(conf.Index))
res, err := copy.Do(context.Background())
if err != nil {
return fmt.Errorf("failed to copy indices: %w", esErrorString(err))
}
table := printer.NewTable(conf).WithSize(2, 0).
WithHeaders("Reindex metrtic", "value")
table.Entries = [][]any{
{"Source indices", conf.SourceIndices},
{"Target index", conf.Index},
{"Batches", *res.Batches},
{"Documents total", *res.Total},
{"Documents created", *res.Created},
{"Documents deleted", *res.Deleted},
{"Documents updated", *res.Updated},
{"Requests/s", *res.RequestsPerSecond},
{"Timed out", *res.TimedOut},
{"Time elapsed", time.Duration(*res.Took) * time.Millisecond},
{"Version conflicts", *res.VersionConflicts},
}
return table.Print()
}

View File

@@ -28,12 +28,29 @@ import (
"codeberg.org/scip/esctl/pkg/cfg" "codeberg.org/scip/esctl/pkg/cfg"
"codeberg.org/scip/esctl/pkg/printer" "codeberg.org/scip/esctl/pkg/printer"
"codeberg.org/scip/mapmap"
"github.com/elastic/go-elasticsearch/v9/typedapi/esdsl" "github.com/elastic/go-elasticsearch/v9/typedapi/esdsl"
"github.com/elastic/go-elasticsearch/v9/typedapi/types" "github.com/elastic/go-elasticsearch/v9/typedapi/types"
) )
// used for completion // used for completion
func IndexTemplateList(conf *cfg.Config) error { func IndexTemplateNames(conf *cfg.Config) ([]string, error) {
res, err := conf.DefaultCluster.ES().Indices.GetIndexTemplate().
Do(context.Background())
if err != nil {
return nil, fmt.Errorf("failed to get index templates: %w", esErrorString(err))
}
names := make([]string, len(res.IndexTemplates))
for idx, tpl := range res.IndexTemplates {
names[idx] = tpl.Name
}
return names, nil
}
func IndexTemplateList(conf *cfg.Config, filter string) error {
res, err := conf.DefaultCluster.ES().Indices.GetIndexTemplate(). res, err := conf.DefaultCluster.ES().Indices.GetIndexTemplate().
Do(context.Background()) Do(context.Background())
if err != nil { if err != nil {
@@ -42,16 +59,20 @@ func IndexTemplateList(conf *cfg.Config) error {
slog.Debug("res", "index templates", res) slog.Debug("res", "index templates", res)
table := printer.NewTable(conf, 5, len(res.IndexTemplates)) table := printer.NewTable(conf).WithSize(5, 0)
table.Addheaders("name", "description", "priority") table.Addheaders("name", "description", "index patterns", "priority")
for idx, tpl := range res.IndexTemplates { tplList := filterIndexTemplates(conf, filter, res.IndexTemplates)
desc, err := json.Marshal(tpl.IndexTemplate.Meta_["description"])
if err != nil { for _, tpl := range tplList {
return fmt.Errorf("failed to unmarshal meta json data: %w", err) desc := strings.TrimPrefix(strings.TrimSuffix(string(tpl.IndexTemplate.Meta_["description"]), `"`), `"`)
var prio int64
if tpl.IndexTemplate.Priority != nil {
prio = *tpl.IndexTemplate.Priority
} }
table.Entries[idx] = []any{tpl.Name, desc, tpl.IndexTemplate.Priority} table.AddRow(tpl.Name, desc, tpl.IndexTemplate.IndexPatterns, prio)
} }
table.Sort() table.Sort()
@@ -59,6 +80,30 @@ func IndexTemplateList(conf *cfg.Config) error {
return table.Print() return table.Print()
} }
// Filter index templates by name, index pattern or hidden flag, using
// mapmap.Slicer
func filterIndexTemplates(conf *cfg.Config, nameFilter string,
templates []types.IndexTemplateItem) []types.IndexTemplateItem {
return mapmap.NewSlicer(templates).MapSliceValuesImmutable(func(tpl types.IndexTemplateItem) bool {
if !conf.Hidden && strings.HasPrefix(tpl.Name, ".") {
return false
}
if nameFilter != "" && !strings.Contains(tpl.Name, nameFilter) {
return false
}
if len(conf.Filter) > 0 {
if !mapmap.NewSlicer(tpl.IndexTemplate.IndexPatterns).
FindSliceInSlice(conf.Filter, strings.Contains) {
return false
}
}
return true
})
}
func IndexTemplateShow(conf *cfg.Config, tplname string) error { func IndexTemplateShow(conf *cfg.Config, tplname string) error {
res, err := conf.DefaultCluster.ES().Indices.GetIndexTemplate(). res, err := conf.DefaultCluster.ES().Indices.GetIndexTemplate().
Name(tplname). Name(tplname).
@@ -75,7 +120,7 @@ func IndexTemplateShow(conf *cfg.Config, tplname string) error {
tpl := res.IndexTemplates[0] tpl := res.IndexTemplates[0]
table := printer.NewTable(conf, 2, 6) table := printer.NewTable(conf).WithSize(2, 6)
table.Addheaders("index template property", "value") table.Addheaders("index template property", "value")
desc, err := json.Marshal(tpl.IndexTemplate.Meta_["description"]) desc, err := json.Marshal(tpl.IndexTemplate.Meta_["description"])
@@ -103,7 +148,7 @@ func IndexTemplateShow(conf *cfg.Config, tplname string) error {
return err return err
} }
table = printer.NewTable(conf, 2, 0) table = printer.NewTable(conf).WithSize(2, 0)
table.Addheaders("index setting property", "value") table.Addheaders("index setting property", "value")
err = getIndexTemplateSettings(conf, tplname, table) err = getIndexTemplateSettings(conf, tplname, table)
@@ -117,7 +162,7 @@ func IndexTemplateShow(conf *cfg.Config, tplname string) error {
return err return err
} }
table = printer.NewTable(conf, 2, 0) table = printer.NewTable(conf).WithSize(2, 0)
table.Addheaders("index field mapping", "type") table.Addheaders("index field mapping", "type")
if tpl.IndexTemplate.Template.Mappings != nil { if tpl.IndexTemplate.Template.Mappings != nil {
@@ -396,7 +441,7 @@ func rolloverAliasIndexTemplate(conf *cfg.Config, name string) error {
return nil return nil
} }
table := printer.NewTable(conf, 4, 0) table := printer.NewTable(conf).WithSize(4, 0)
table.Addheaders("rollover alias", "status", "ack", "new index") table.Addheaders("rollover alias", "status", "ack", "new index")
// apply rollover to all matching aliases, if any // apply rollover to all matching aliases, if any

View File

@@ -35,7 +35,7 @@ func LicenseShow(conf *cfg.Config) error {
slog.Debug("license show", "license", res) slog.Debug("license show", "license", res)
table := printer.NewTableEmpty(conf).WithHeaders("license setting", "value") table := printer.NewTable(conf).WithHeaders("license setting", "value")
lic := res.License lic := res.License

View File

@@ -26,7 +26,6 @@ import (
"codeberg.org/scip/esctl/pkg/cfg" "codeberg.org/scip/esctl/pkg/cfg"
"codeberg.org/scip/esctl/pkg/printer" "codeberg.org/scip/esctl/pkg/printer"
"github.com/dustin/go-humanize"
) )
func NodeList(conf *cfg.Config) error { func NodeList(conf *cfg.Config) error {
@@ -38,11 +37,11 @@ func NodeList(conf *cfg.Config) error {
slog.Debug("ES result", "nodes", nodes) slog.Debug("ES result", "nodes", nodes)
table := printer.NewTable(conf, 7, len(nodes)) table := printer.NewTable(conf).WithHeaders(
table.Addheaders("name", "ip", "load1m", "load5m", "load15m", "ram %", "heap %") "name", "ip", "load1m", "load5m", "load15m", "ram %", "heap %")
for idx, node := range nodes { for _, node := range nodes {
table.Entries[idx] = []any{ table.AddRow(
*node.Name, *node.Name,
*node.Ip, *node.Ip,
*node.Load1M, *node.Load1M,
@@ -50,7 +49,7 @@ func NodeList(conf *cfg.Config) error {
*node.Load15M, *node.Load15M,
node.RamPercent, node.RamPercent,
node.HeapPercent, node.HeapPercent,
} )
} }
table.Sort() table.Sort()
@@ -105,7 +104,7 @@ func NodeShow(conf *cfg.Config, nodename string) error {
for id, info := range res.Nodes { for id, info := range res.Nodes {
stat := stats.Nodes[id] stat := stats.Nodes[id]
table := printer.NewTable(conf, 2, 0) table := printer.NewTable(conf).WithSize(2, 0)
table.Addheaders(nodename+" property", "value") table.Addheaders(nodename+" property", "value")
roles := make([]string, len(info.Roles)) roles := make([]string, len(info.Roles))
@@ -114,18 +113,27 @@ func NodeShow(conf *cfg.Config, nodename string) error {
} }
k8snode := info.Attributes["k8s_node_name"] k8snode := info.Attributes["k8s_node_name"]
rank := "none"
adsel, exists := stat.AdaptiveSelection[id]
if exists {
rank = *adsel.Rank
}
table.Entries = [][]any{ table.Entries = [][]any{
{"Id", id}, {"Id", id},
{"Name", nodename}, {"Name", nodename},
{"Kubernetes node", k8snode}, {"Kubernetes node", k8snode},
{"Ip address", info.Ip}, {"Ip address", info.Ip},
{"Node rank", *stat.AdaptiveSelection[id].Rank}, {"Node rank", rank},
{"JVM", info.Jvm.VmName + " " + info.Jvm.Version}, {"JVM", info.Jvm.VmName + " " + info.Jvm.Version},
{"JVM Started", time.UnixMilli(info.Jvm.StartTimeInMillis)}, {"JVM Started", time.UnixMilli(info.Jvm.StartTimeInMillis)},
{"OS", info.Os.PrettyName + " " + info.Os.Version}, {"OS", info.Os.PrettyName + " " + info.Os.Version},
{"Node roles", roles}, {"Node roles", roles},
{"Node version", info.Version}, {"Node version", info.Version},
// FIXME: not implemented upstream
// see: https://github.com/elastic/go-elasticsearch/issues/1526
// {"Allocated shards", stat.Allocations.XXX},
{"HTTP clients", *stat.Http.CurrentOpen}, {"HTTP clients", *stat.Http.CurrentOpen},
{"CPUs", *info.Os.AllocatedProcessors}, {"CPUs", *info.Os.AllocatedProcessors},
{"Load 15m/5m/1m", fmt.Sprintf("%.2f/%.2f/%.2f", {"Load 15m/5m/1m", fmt.Sprintf("%.2f/%.2f/%.2f",
@@ -134,18 +142,41 @@ func NodeShow(conf *cfg.Config, nodename string) error {
stat.Os.Cpu.LoadAverage["1m"], stat.Os.Cpu.LoadAverage["1m"],
)}, )},
{"Open FD's", *stat.Process.OpenFileDescriptors}, {"Open FD's", *stat.Process.OpenFileDescriptors},
{"Response time avg", fmt.Sprintf("%dns", *stat.AdaptiveSelection[id].AvgResponseTimeNs)}, {"HTTP sesssions current/total", fmt.Sprintf("%d/%d",
*stat.Http.CurrentOpen,
*stat.Http.TotalOpened,
)},
{"Traffic rx/tx",
printer.ByteString(*stat.Transport.RxSizeInBytes) + " / " + printer.ByteString(*stat.Transport.TxSizeInBytes)},
{"Response time avg", time.Duration(*stat.AdaptiveSelection[id].AvgResponseTimeNs)},
{"Memory usage (used/avail)", {"Memory usage (used/avail)",
humanize.Bytes(uint64(*stat.Os.Mem.UsedInBytes)) + " / " + humanize.Bytes(uint64(*stat.Os.Mem.TotalInBytes))}, printer.ByteString(*stat.Os.Mem.UsedInBytes) + " / " + printer.ByteString(*stat.Os.Mem.TotalInBytes)},
{"Search queries current/total", fmt.Sprintf("%d/%d",
stat.Indices.Search.QueryCurrent,
stat.Indices.Search.QueryTotal,
)},
{"Search efficiency", stat.Indices.Search.QueryTimeInMillis / stat.Indices.Search.QueryTotal},
{"Docs count", stat.Indices.Docs.Count},
{"Merges current/total", fmt.Sprintf("%d/%d",
stat.Indices.Merges.Current,
stat.Indices.Merges.Total,
)},
{"Merge docs count current/total", fmt.Sprintf("%d/%d",
stat.Indices.Merges.CurrentDocs,
stat.Indices.Merges.TotalDocs,
)},
{"Merge size current/total", fmt.Sprintf("%s/%s",
printer.ByteString(stat.Indices.Merges.CurrentSizeInBytes),
printer.ByteString(stat.Indices.Merges.TotalSizeInBytes),
)},
{"CircuitBreaker trip count", *stat.Breakers["fielddata"].Tripped},
} }
if len(stat.Fs.Data) > 0 { if len(stat.Fs.Data) > 0 {
fs := stat.Fs.Data[0] fs := stat.Fs.Data[0]
table.Entries = append(table.Entries, [][]any{ table.AddRow("Storage usage (used/avail)",
{"Storage usage (used/avail)", printer.ByteString(*fs.AvailableInBytes)+" / "+printer.ByteString(*fs.TotalInBytes))
humanize.Bytes(uint64(*fs.AvailableInBytes)) + " / " + humanize.Bytes(uint64(*fs.TotalInBytes))}, table.AddRow("Storage mount", *fs.Mount)
{"Storage mount", *fs.Mount},
}...)
} }
if err := table.Print(); err != nil { if err := table.Print(); err != nil {
@@ -166,7 +197,7 @@ func NodeClients(conf *cfg.Config, nodename string) error {
slog.Debug("ES result", "stat", stats) slog.Debug("ES result", "stat", stats)
table := printer.NewTable(conf, 5, 0) table := printer.NewTable(conf).WithSize(5, 0)
table.Addheaders("agent", "id", "when", "from host", "url") table.Addheaders("agent", "id", "when", "from host", "url")
for _, stat := range stats.Nodes { for _, stat := range stats.Nodes {
@@ -204,3 +235,72 @@ func NodeClients(conf *cfg.Config, nodename string) error {
return table.Print() return table.Print()
} }
func NodeUsage(conf *cfg.Config, nodeid string) error {
usage := conf.DefaultCluster.ES().Nodes.Usage()
if nodeid != "" {
usage.NodeId(nodeid)
}
stats, err := usage.Do(context.Background())
if err != nil {
return fmt.Errorf("failed to get node usage: %w", esErrorString(err))
}
slog.Debug("ES result", "usage", stats)
table := printer.NewTable(conf).WithHeaders(
"node",
"bulk",
"doc get",
"doc mget",
"doc update",
"index doc",
"index stats",
"search",
"msearch",
"open pit",
)
for id, actions := range stats.Nodes {
node, err := getNodeName(conf, id)
if err != nil {
return err
}
stat := actions.RestActions
table.AddRow(
node,
stat["bulk_action"],
stat["document_get_action"],
stat["document_mget_action"],
stat["document_update_action"],
stat["document_index_action"],
stat["indices_stats_action"],
stat["search_action"],
stat["msearch_action"],
stat["open_point_in_time"],
)
}
return table.Print()
}
func getNodeName(conf *cfg.Config, id string) (string, error) {
res, err := conf.DefaultCluster.ES().Nodes.Info().
Metric("os").
Do(context.Background())
if err != nil {
return "", fmt.Errorf("failed to get node info: %w", esErrorString(err))
}
for nodeid, node := range res.Nodes {
if id == nodeid {
return node.Name, nil
}
}
return "", nil
}

View File

@@ -19,7 +19,6 @@ package es
import ( import (
"context" "context"
"fmt" "fmt"
"sync"
"codeberg.org/scip/esctl/pkg/cfg" "codeberg.org/scip/esctl/pkg/cfg"
"github.com/elastic/go-elasticsearch/v9" "github.com/elastic/go-elasticsearch/v9"
@@ -61,14 +60,7 @@ type apiResponse struct {
which int which int
} }
func getApiData( func getApiData(conf *cfg.Config, es *elasticsearch.TypedClient, reschan chan apiResponse, which string) {
conf *cfg.Config,
es *elasticsearch.TypedClient,
wg *sync.WaitGroup,
reschan chan apiResponse,
which string) {
defer wg.Done()
apiRes := apiResponse{} apiRes := apiResponse{}
var arerr error var arerr error

View File

@@ -20,6 +20,8 @@ import (
"context" "context"
"fmt" "fmt"
"log/slog" "log/slog"
"maps"
"slices"
"codeberg.org/scip/esctl/pkg/cfg" "codeberg.org/scip/esctl/pkg/cfg"
"codeberg.org/scip/esctl/pkg/printer" "codeberg.org/scip/esctl/pkg/printer"
@@ -33,15 +35,7 @@ func RoleNames(conf *cfg.Config) ([]string, error) {
return nil, fmt.Errorf("failed to get roles: %w", esErrorString(err)) return nil, fmt.Errorf("failed to get roles: %w", esErrorString(err))
} }
roles := make([]string, len(res)) return slices.Collect(maps.Keys(res)), nil
idx := 0
for name := range res {
roles[idx] = name
idx++
}
return roles, nil
} }
func RoleList(conf *cfg.Config) error { func RoleList(conf *cfg.Config) error {
@@ -53,7 +47,7 @@ func RoleList(conf *cfg.Config) error {
slog.Debug("ES result", "roles", res) slog.Debug("ES result", "roles", res)
table := printer.NewTable(conf, 3, len(res)) table := printer.NewTable(conf).WithSize(3, len(res))
table.Addheaders("role", "index roles", "cluster roles") table.Addheaders("role", "index roles", "cluster roles")
idx := 0 idx := 0
@@ -132,7 +126,7 @@ func roleRemoteClusters(conf *cfg.Config, role types.Role) error {
return nil return nil
} }
table := printer.NewTable(conf, 2, len(role.RemoteCluster)) table := printer.NewTable(conf).WithSize(2, len(role.RemoteCluster))
table.Addheaders("remote cluster", "privilege") table.Addheaders("remote cluster", "privilege")
idx := 0 idx := 0
@@ -158,7 +152,7 @@ func roleClusters(conf *cfg.Config, role types.Role) error {
return nil return nil
} }
table := printer.NewTable(conf, 1, len(role.Cluster)) table := printer.NewTable(conf).WithSize(1, len(role.Cluster))
table.Addheaders("cluster rights") table.Addheaders("cluster rights")
idx := 0 idx := 0
@@ -177,7 +171,7 @@ func roleRemoteIndices(conf *cfg.Config, role types.Role) error {
return nil return nil
} }
table := printer.NewTable(conf, 3, len(role.RemoteIndices)) table := printer.NewTable(conf).WithSize(3, len(role.RemoteIndices))
table.Addheaders("remote index names", "index permissions", "allow restricted") table.Addheaders("remote index names", "index permissions", "allow restricted")
idx := 0 idx := 0
@@ -203,7 +197,7 @@ func roleIndices(conf *cfg.Config, role types.Role) error {
return nil return nil
} }
table := printer.NewTable(conf, 3, len(role.Indices)) table := printer.NewTable(conf).WithSize(3, len(role.Indices))
table.Addheaders("index names", "index permissions", "allow restricted") table.Addheaders("index names", "index permissions", "allow restricted")
idx := 0 idx := 0
@@ -229,7 +223,7 @@ func roleApplications(conf *cfg.Config, role types.Role) error {
return nil return nil
} }
table := printer.NewTable(conf, 3, len(role.Applications)) table := printer.NewTable(conf).WithSize(3, len(role.Applications))
table.Addheaders("application", "privileges", "resources") table.Addheaders("application", "privileges", "resources")
idx := 0 idx := 0

View File

@@ -127,7 +127,7 @@ func getCsvRecord(conf *cfg.Config, csvfile, rolename string) (*Record, error) {
}() }()
scanner := bufio.NewScanner(fd) scanner := bufio.NewScanner(fd)
record := Record{role: rolename} record := new(Record{role: rolename})
for scanner.Scan() { for scanner.Scan() {
line := strings.TrimSpace(scanner.Text()) line := strings.TrimSpace(scanner.Text())
@@ -150,7 +150,7 @@ func getCsvRecord(conf *cfg.Config, csvfile, rolename string) (*Record, error) {
} }
} }
return &record, nil return record, nil
} }
func diffRoles(conf *cfg.Config, records map[string]Record, res getrole.Response) []Register { func diffRoles(conf *cfg.Config, records map[string]Record, res getrole.Response) []Register {
@@ -216,7 +216,7 @@ func RoleDiff(conf *cfg.Config, csvfile, role string) error {
rows := diffRoles(conf, records, res) rows := diffRoles(conf, records, res)
table := printer.NewTable(conf, 3, len(rows)) table := printer.NewTable(conf).WithSize(3, len(rows))
table.Addheaders("role", "is deployed", "is defined") table.Addheaders("role", "is deployed", "is defined")
for idx, row := range rows { for idx, row := range rows {

View File

@@ -33,7 +33,7 @@ func RolloverConditions(conf *cfg.Config) types.RolloverConditionsVariant {
} }
if conf.MaxDocs > 0 { if conf.MaxDocs > 0 {
cond.MaxDocs(int64(conf.MaxDocs)) cond.MaxDocs(conf.MaxDocs)
} }
if conf.MaxShardSize > 0 { if conf.MaxShardSize > 0 {

View File

@@ -56,7 +56,7 @@ func Search(conf *cfg.Config, queries []string) error {
return err return err
} }
req := &search.Request{Query: queryCaster} req := new(search.Request{Query: queryCaster})
searchEs.Request(req) searchEs.Request(req)
@@ -128,7 +128,7 @@ func validateSearch(conf *cfg.Config, queries []string) error {
return err return err
} }
req := &validatequery.Request{Query: queryCaster} req := new(validatequery.Request{Query: queryCaster})
validate.Request(req) validate.Request(req)

View File

@@ -80,7 +80,7 @@ func NewFilter(query string) (*filter, error) {
return nil, errors.New("search queries must be in the form field<sep>pattern where <sep> must be one of: = or !=") return nil, errors.New("search queries must be in the form field<sep>pattern where <sep> must be one of: = or !=")
} }
flt := &filter{term: part[0], filter: part[1], criteria: criteria} flt := new(filter{term: part[0], filter: part[1], criteria: criteria})
if strings.Contains(part[0], ",") { if strings.Contains(part[0], ",") {
// a MultiMatchQuery, match across multiple fields at once // a MultiMatchQuery, match across multiple fields at once

79
pkg/es/searchql.go Normal file
View File

@@ -0,0 +1,79 @@
/*
Copyright © 2026 Thomas von Dein
This program is free software: you can redistribute it and/or modify
it under the terms of the GNU General Public License as published by
the Free Software Foundation, either version 3 of the License, or
(at your option) any later version.
This program is distributed in the hope that it will be useful,
but WITHOUT ANY WARRANTY; without even the implied warranty of
MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
GNU General Public License for more details.
You should have received a copy of the GNU General Public License
along with this program. If not, see <http://www.gnu.org/licenses/>.
*/
package es
import (
"context"
"encoding/json"
"fmt"
"log/slog"
"codeberg.org/scip/esctl/pkg/cfg"
"codeberg.org/scip/esctl/pkg/printer"
)
type searchQlColumn struct {
Name string `json:"name"`
Type string `json:"type"`
}
type searchQlResult struct {
Partial bool `json:"is_partial"`
Docs int `json:"documents_found"`
Columns []searchQlColumn `json:"columns"`
Values [][]any `json:"values"`
}
func SearchQL(conf *cfg.Config, querystring string) error {
query := conf.DefaultCluster.ES().Esql.Query().
Query(querystring).
DropNullColumns(true)
res, err := query.Do(context.Background())
if err != nil {
return fmt.Errorf("failed to run esql query: %w", esErrorString(err))
}
if conf.Debug {
fmt.Println(string(res))
}
qlResult := searchQlResult{}
if err = json.Unmarshal(res, &qlResult); err != nil {
return fmt.Errorf("failed to unmarshal esql result: %w", err)
}
slog.Debug("searchql result", "qlres", qlResult)
if qlResult.Docs == 0 {
return nil
}
table := printer.NewTable(conf).WithSize(len(qlResult.Columns), len(qlResult.Values))
headers := make([]string, len(qlResult.Columns))
for idx, col := range qlResult.Columns {
headers[idx] = col.Name
}
table.Addheaders(headers...)
table.Entries = qlResult.Values
return table.Print()
}

View File

@@ -100,7 +100,7 @@ func printShards(conf *cfg.Config, shardlist shards.Response) error {
headers = append(headers, "node", "ip") headers = append(headers, "node", "ip")
} }
table := printer.NewTable(conf, len(headers), len(shardlist)) table := printer.NewTable(conf).WithSize(len(headers), len(shardlist))
table.Addheaders(headers...) table.Addheaders(headers...)
for idx, shard := range shardlist { for idx, shard := range shardlist {
@@ -163,28 +163,35 @@ func ShardAllocation(conf *cfg.Config, index string) error {
slog.Debug("ES result", "explain", res) slog.Debug("ES result", "explain", res)
currentNode := res.CurrentNode table := printer.NewTable(conf).WithSize(2, 10)
table := printer.NewTable(conf, 2, 10)
table.Addheaders("shard allocation setting", "value") table.Addheaders("shard allocation setting", "value")
roles := make([]string, len(currentNode.Roles))
for idx, role := range currentNode.Roles {
roles[idx] = role.Name
}
table.Entries = [][]any{ table.Entries = [][]any{
{"Index", index}, {"Index", index},
{"Current state", res.CurrentState}, {"Current state", res.CurrentState},
{"Can rebalance cluster", res.CanRebalanceCluster.Name},
{"Can rebalance to another node", res.CanRebalanceToOtherNode.Name},
{"Can remain on current node", res.CanRemainOnCurrentNode.Name},
}
if res.CurrentNode != nil {
currentNode := res.CurrentNode
roles := make([]string, len(currentNode.Roles))
for idx, role := range currentNode.Roles {
roles[idx] = role.Name
}
table.Entries = append(table.Entries, [][]any{
{"Current node", currentNode.Name}, {"Current node", currentNode.Name},
{"Current k8s node", currentNode.Attributes["k8s_node_name"]}, {"Current k8s node", currentNode.Attributes["k8s_node_name"]},
{"Current node address", currentNode.TransportAddress}, {"Current node address", currentNode.TransportAddress},
{"Current node id", currentNode.Id}, {"Current node id", currentNode.Id},
{"Current node weight", currentNode.WeightRanking}, {"Current node weight", currentNode.WeightRanking},
{"Current node roles", roles}, {"Current node roles", roles},
{"Can rebalance cluster", res.CanRebalanceCluster.Name}, }...)
{"Can rebalance to another node", res.CanRebalanceToOtherNode.Name}, } else {
{"Can remain on current node", res.CanRemainOnCurrentNode.Name}, table.AddRow("Current node", "not currently assigned to any node")
} }
if res.CurrentState == "unassigned" { if res.CurrentState == "unassigned" {

View File

@@ -42,17 +42,17 @@ type Snapshot struct {
} }
func SnapshotList(conf *cfg.Config) error { func SnapshotList(conf *cfg.Config) error {
// get partial indicies // get partial indices
ires, err := conf.DefaultCluster.ES().Cat.Indices().Do(context.Background()) ires, err := conf.DefaultCluster.ES().Cat.Indices().Do(context.Background())
if err != nil { if err != nil {
return fmt.Errorf("failed to get indicies: %w", esErrorString(err)) return fmt.Errorf("failed to get indices: %w", esErrorString(err))
} }
indicies := map[string]int{} indices := map[string]int{}
for _, index := range ires { for _, index := range ires {
name := strings.ReplaceAll(*index.Index, "partial-", "") name := strings.ReplaceAll(*index.Index, "partial-", "")
indicies[name] = 1 indices[name] = 1
} }
// get snapshots // get snapshots
@@ -61,20 +61,20 @@ func SnapshotList(conf *cfg.Config) error {
return fmt.Errorf("failed to get snapshots: %w", esErrorString(err)) return fmt.Errorf("failed to get snapshots: %w", esErrorString(err))
} }
slog.Debug("ES result", "indicies", sres) slog.Debug("ES result", "indices", sres)
snapshots := []*Snapshot{} // original snapshot names snapshots := []*Snapshot{} // original snapshot names
for _, snapshot := range sres { for _, snapshot := range sres {
snap := &Snapshot{ snap := new(Snapshot{
Name: *snapshot.Id, Name: *snapshot.Id,
Status: *snapshot.Status, Status: *snapshot.Status,
Start: fmt.Sprintf("%s", snapshot.StartTime), Start: fmt.Sprintf("%s", snapshot.StartTime),
Forindex: indexFromSnapshot(*snapshot.Id), Forindex: indexFromSnapshot(*snapshot.Id),
Orphaned: "no", Orphaned: "no",
} })
_, exists := indicies[snap.Forindex] _, exists := indices[snap.Forindex]
if !exists { if !exists {
snap.Orphaned = "orphaned" snap.Orphaned = "orphaned"
} }
@@ -84,7 +84,7 @@ func SnapshotList(conf *cfg.Config) error {
} }
} }
table := printer.NewTable(conf, 5, len(snapshots)) table := printer.NewTable(conf).WithSize(5, len(snapshots))
table.Addheaders("name", "index", "start", "orphaned", "status") table.Addheaders("name", "index", "start", "orphaned", "status")
for idx, snap := range snapshots { for idx, snap := range snapshots {
@@ -118,7 +118,7 @@ func SnapshotShow(conf *cfg.Config, snapshot string) error {
return errors.New("no snapshot retrieved") return errors.New("no snapshot retrieved")
} }
table := printer.NewTable(conf, 2, 17) table := printer.NewTable(conf).WithSize(2, 17)
table.Addheaders("snapshot property", "value") table.Addheaders("snapshot property", "value")
snap := res.Snapshots[0] snap := res.Snapshots[0]

View File

@@ -36,7 +36,7 @@ func TaskList(conf *cfg.Config) error {
slog.Debug("res", "tasks", res) slog.Debug("res", "tasks", res)
table := printer.NewTable(conf, 8, len(res)) table := printer.NewTable(conf).WithSize(6, len(res))
table.Addheaders("task id", "action", "start time", "run time", "node", "type") table.Addheaders("task id", "action", "start time", "run time", "node", "type")
for idx, task := range res { for idx, task := range res {

View File

@@ -29,13 +29,13 @@ import (
const LevelNotice = slog.Level(2) const LevelNotice = slog.Level(2)
func Init(conf *cfg.Config) { func Init(conf *cfg.Config) {
logLevel := &slog.LevelVar{} logLevel := new(slog.LevelVar{})
opts := &yadu.Options{ opts := new(yadu.Options{
Level: logLevel, Level: logLevel,
AddSource: true, AddSource: true,
NoColor: !isatty.IsTerminal(os.Stdout.Fd()), NoColor: !isatty.IsTerminal(os.Stdout.Fd()),
} })
buildInfo, _ := debug.ReadBuildInfo() buildInfo, _ := debug.ReadBuildInfo()

View File

@@ -22,10 +22,16 @@ type ByteSize struct {
size uint64 size uint64
} }
func Bytes(size int64) ByteSize {
return ByteSize{size: uint64(size)}
}
func (b *ByteSize) String() string { func (b *ByteSize) String() string {
return humanize.Bytes(b.size) return humanize.Bytes(b.size)
} }
func Bytes(size int64) *ByteSize {
return new(ByteSize{size: uint64(size)})
}
func ByteString(size int64) string {
b := ByteSize{size: uint64(size)}
return b.String()
}

View File

@@ -1,3 +1,5 @@
package printer
/* /*
Copyright © 2026 Thomas von Dein Copyright © 2026 Thomas von Dein
@@ -14,7 +16,6 @@ GNU General Public License for more details.
You should have received a copy of the GNU General Public License You should have received a copy of the GNU General Public License
along with this program. If not, see <http://www.gnu.org/licenses/>. along with this program. If not, see <http://www.gnu.org/licenses/>.
*/ */
package printer
import ( import (
"fmt" "fmt"
@@ -25,7 +26,9 @@ import (
"github.com/elastic/go-elasticsearch/v9/typedapi/types" "github.com/elastic/go-elasticsearch/v9/typedapi/types"
) )
func any2string(in any) string { // stringer returns the string representation of different types of
// values
func stringer(in any) string {
//nolint:gocritic //nolint:gocritic
switch val := in.(type) { switch val := in.(type) {
case string: case string:
@@ -38,62 +41,42 @@ func any2string(in any) string {
return strconv.Itoa(val) return strconv.Itoa(val)
case float64: case float64:
return fmt.Sprintf("%.2f", val) return fmt.Sprintf("%.2f", val)
case float32:
return fmt.Sprintf("%.2f", val)
case []string: case []string:
return strings.Join(val, ",") return strings.Join(val, ",")
case ByteSize: case ByteSize:
return val.String() return val.String()
case *ByteSize:
return val.String()
case time.Time: case time.Time:
return val.Format("2006-01-02 15:04:05") return val.Format("2006-01-02 15:04:05")
case time.Duration:
return val.String()
case []byte: case []byte:
return string(val) return string(val)
case nil: case nil:
return "null" return "null"
case types.Percentage, types.DateTime: case *types.Float64:
return val.(string) // ignore err here, because types.Float64.MarshalJSON() never returns one
} f, _ := val.MarshalJSON()
return string(f)
return "" default:
} // Caution: this may cause a panic if the [unknown] type does
// not implement fmt.Stringer. In this case add another case
func (data *Table) preprocessRows() { // to the type switch for it above.
if data.processed { return in.(fmt.Stringer).String()
// only do it once
return
}
// convert entries to strings
data.rows = make([][]string, len(data.Entries))
for rowidx, entries := range data.Entries {
data.rows[rowidx] = make([]string, len(entries))
for colidx, entry := range data.Entries[rowidx] {
data.rows[rowidx][colidx] = any2string(entry)
} }
} }
// determine header lenght's // visibleLen returns the length of a string but only visible chars,
for idx, head := range data.Headers { // w/o ansi color escapes
data.lenHeaders[idx] = visibleLen(head) func visibleLen(word string) int {
if !strings.Contains(word, "\x1b") {
// no ansi escape in there, use faster method
return len(word)
} }
// determine max width per column // contains escapes, need to clean up before counting
for _, entries := range data.rows { return len(ansiCtrlSeq.ReplaceAllLiteralString(word, ""))
currentWidth := 0
for idx, entry := range entries {
length := visibleLen(entry)
if data.lenHeaders[idx] < length {
if length > currentWidth+data.maxwidth {
data.lenHeaders[idx] = data.maxwidth - currentWidth
} else {
data.lenHeaders[idx] = length
}
}
currentWidth += data.lenHeaders[idx]
}
}
data.processed = true
} }

View File

@@ -18,8 +18,11 @@ along with this program. If not, see <http://www.gnu.org/licenses/>.
package printer package printer
import ( import (
"os"
"codeberg.org/scip/esctl/pkg/cfg" "codeberg.org/scip/esctl/pkg/cfg"
"github.com/fatih/color" "github.com/fatih/color"
"github.com/mattn/go-isatty"
) )
var ( var (
@@ -31,6 +34,10 @@ var (
) )
func Colorize(conf *cfg.Config, col, what string) string { func Colorize(conf *cfg.Config, col, what string) string {
if !isatty.IsTerminal(os.Stdout.Fd()) {
return what
}
switch conf.Output { switch conf.Output {
case "json", "yaml": case "json", "yaml":
return what return what
@@ -49,3 +56,11 @@ func Colorize(conf *cfg.Config, col, what string) string {
return what return what
} }
func Bold(what string) string {
if !isatty.IsTerminal(os.Stdout.Fd()) {
return what
}
return bold(what)
}

55
pkg/printer/metrics.go Normal file
View File

@@ -0,0 +1,55 @@
/*
Copyright © 2026 Thomas von Dein
This program is free software: you can redistribute it and/or modify
it under the terms of the GNU General Public License as published by
the Free Software Foundation, either version 3 of the License, or
(at your option) any later version.
This program is distributed in the hope that it will be useful,
but WITHOUT ANY WARRANTY; without even the implied warranty of
MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
GNU General Public License for more details.
You should have received a copy of the GNU General Public License
along with this program. If not, see <http://www.gnu.org/licenses/>.
*/
package printer
import (
"fmt"
"runtime/metrics"
"strconv"
)
func printGoRoutineMetrics() {
fmt.Println("\nGoroutine metrics:")
printMetric("/sched/goroutines-created:goroutines", "Created")
printMetric("/sched/goroutines:goroutines", "Live")
printMetric("/sched/goroutines/not-in-go:goroutines", "Syscall/CGO")
printMetric("/sched/goroutines/runnable:goroutines", "Runnable")
printMetric("/sched/goroutines/running:goroutines", "Running")
printMetric("/sched/goroutines/waiting:goroutines", "Waiting")
fmt.Println("Thread metrics:")
printMetric("/sched/gomaxprocs:threads", "Max")
printMetric("/sched/threads/total:threads", "Live")
}
func printMetric(name string, descr string) {
sample := []metrics.Sample{{Name: name}}
metrics.Read(sample)
var val string
switch sample[0].Value.Kind() {
case metrics.KindFloat64, metrics.KindFloat64Histogram:
val = fmt.Sprintf("%.2f", sample[0].Value.Float64())
case metrics.KindUint64:
val = strconv.FormatUint(sample[0].Value.Uint64(), 10)
case metrics.KindBad:
val = "n/a"
}
fmt.Printf(" %s: %v\n", descr, val)
}

View File

@@ -25,13 +25,13 @@ import (
"strings" "strings"
"codeberg.org/scip/esctl/pkg/cfg" "codeberg.org/scip/esctl/pkg/cfg"
"github.com/seeruk/go-wordwrap"
"gopkg.in/yaml.v3" "gopkg.in/yaml.v3"
) )
type Table struct { type Table struct {
Mode string // tsv, json, yaml Mode string // tsv, json, yaml
Headers []string Headers []string
RawHeaders []string
Entries [][]any Entries [][]any
rows [][]string // representation used for printing rows [][]string // representation used for printing
@@ -39,47 +39,62 @@ type Table struct {
lenHeaders []int lenHeaders []int
alignInts bool alignInts bool
maxwidth int maxwidth int
debugGoRoutines bool
} }
func NewTable(conf *cfg.Config, columns, rows int) *Table { func NewTable(conf *cfg.Config) *Table {
table := Table{Mode: conf.Output, maxwidth: cfg.GetTermWidth()} return new(Table{
Mode: conf.Output,
maxwidth: cfg.GetTermWidth(),
debugGoRoutines: conf.DebugGoRoutines,
Headers: []string{},
RawHeaders: []string{},
lenHeaders: []int{},
alignInts: conf.AlignInts,
})
}
func (table *Table) WithSize(columns, rows int) *Table {
if columns > 0 {
table.Headers = make([]string, columns) table.Headers = make([]string, columns)
table.Entries = make([][]any, rows) table.RawHeaders = make([]string, columns)
table.lenHeaders = make([]int, columns) table.lenHeaders = make([]int, columns)
table.alignInts = conf.AlignInts
return &table
} }
func NewTableEmpty(conf *cfg.Config) *Table { if rows > 0 {
table := Table{Mode: conf.Output, maxwidth: cfg.GetTermWidth()} table.Entries = make([][]any, rows)
table.alignInts = conf.AlignInts }
return &table return table
} }
func (table *Table) WithHeaders(headers ...string) *Table { func (table *Table) WithHeaders(headers ...string) *Table {
count := len(headers) count := len(headers)
table.Entries = [][]any{} table.WithSize(count, 0).Addheaders(headers...)
table.lenHeaders = make([]int, count)
table.Headers = make([]string, count)
table.Addheaders(headers...)
return table return table
} }
func (table *Table) Print() error { func (table *Table) Print() error {
var err error
switch table.Mode { switch table.Mode {
case "json": case "json":
return table.PrintJSON() err = table.PrintJSON()
case "yaml": case "yaml":
return table.PrintYAML() err = table.PrintYAML()
case "csv":
err = table.PrintCSV()
default: default:
return table.PrintTSV() err = table.PrintTSV()
} }
if table.debugGoRoutines {
printGoRoutineMetrics()
}
return err
} }
var ( var (
@@ -134,23 +149,17 @@ func (table *Table) PrintTSV() error {
for _, entries := range table.rows { for _, entries := range table.rows {
currentWidth := 0 currentWidth := 0
columns := len(entries)
for idx, entry := range entries { for idx, entry := range entries {
length := visibleLen(entry) length := visibleLen(entry)
if length+currentWidth > table.maxwidth && table.maxwidth-currentWidth > 1 { if length+currentWidth > table.maxwidth &&
// text is too wide to be put into one line, wrap it table.maxwidth-currentWidth > 1 &&
wrapper := wordwrap.Wrapper(table.maxwidth-currentWidth, false) idx == columns-1 {
wrapped := wrapper(entry) // text is too wide to be put into one line, and
// it's the last cell, so wrap it
// and indent it entry = wrap(table.maxwidth-currentWidth, currentWidth+2, entry)
for idx, line := range strings.Split(wrapped, "\n") {
if idx == 0 {
entry = line
} else {
entry += "\n " + strings.Repeat(" ", currentWidth) + line
}
}
} }
currentWidth += table.lenHeaders[idx] currentWidth += table.lenHeaders[idx]
@@ -178,6 +187,70 @@ func (table *Table) PrintTSV() error {
return nil return nil
} }
// Wrap a text into multiple lines, first line is not indented, all
// further lines will be indented. Used within Print() to print large
// cell text.
func wrap(width, indent int, text string) string {
if len(text) <= width {
return text
}
wrapped := ""
line := ""
for word := range strings.FieldsSeq(text) {
if len(line)+len(word)+1 <= width {
// appending word to current line doesn't exceed width
if line != "" {
line += " "
}
line += word
} else {
// it exceeds it, so we need to wrap
if wrapped == "" {
// beginning of output, no indenting here
wrapped = line + "\n"
} else {
// we're in the middle of the text, so add the indent
wrapped += strings.Repeat(" ", indent) + line + "\n"
}
// remember the current word for the next round
line = word
}
}
if line != "" {
// last line, no newline needed here
wrapped += strings.Repeat(" ", indent) + line
}
return wrapped
}
func (table *Table) PrintCSV() error {
table.preprocessRows()
fmt.Println(strings.Join(table.RawHeaders, ","))
for _, entries := range table.rows {
row := make([]string, len(entries))
for idx, entry := range entries {
if strings.Contains(entry, " ") || strings.Contains(entry, ",") {
row[idx] = `"` + entry + `"`
} else {
row[idx] = entry
}
}
fmt.Println(strings.Join(row, ","))
}
return nil
}
func (table *Table) Sort() { func (table *Table) Sort() {
// sanity checks // sanity checks
if len(table.Entries) == 0 { if len(table.Entries) == 0 {
@@ -197,8 +270,10 @@ func (table *Table) Addheaders(headers ...string) {
case "json", "yaml": case "json", "yaml":
table.Headers[idx] = strings.ReplaceAll(strings.ToLower(header), " ", "_") table.Headers[idx] = strings.ReplaceAll(strings.ToLower(header), " ", "_")
default: default:
table.Headers[idx] = bold(strings.ReplaceAll(strings.ToUpper(header), " ", "-")) table.Headers[idx] = Bold(strings.ReplaceAll(strings.ToUpper(header), " ", "-"))
} }
table.RawHeaders[idx] = header
} }
} }
@@ -206,6 +281,22 @@ func (table *Table) AddRow(fields ...any) {
table.Entries = append(table.Entries, fields) table.Entries = append(table.Entries, fields)
} }
func (table *Table) AddRowLate(fields ...any) {
table.AddRow(fields)
if !table.processed {
return
}
row := make([]string, len(fields))
for idx, field := range fields {
row[idx] = stringer(field)
}
table.rows = append(table.rows, row)
}
// needed for json and yaml output // needed for json and yaml output
func (table *Table) toMap() []map[string]any { func (table *Table) toMap() []map[string]any {
raw := make([]map[string]any, len(table.Entries)) raw := make([]map[string]any, len(table.Entries))
@@ -220,9 +311,47 @@ func (table *Table) toMap() []map[string]any {
return raw return raw
} }
// return the length of a string but only visible chars, w/o ansi color escapes func (data *Table) preprocessRows() {
func visibleLen(word string) int { if data.processed {
return len(ansiCtrlSeq.ReplaceAllLiteralString(word, "")) // only do it once
return
}
// convert entries to strings
data.rows = make([][]string, len(data.Entries))
for rowidx, entries := range data.Entries {
data.rows[rowidx] = make([]string, len(entries))
for colidx, entry := range data.Entries[rowidx] {
data.rows[rowidx][colidx] = stringer(entry)
}
}
// determine header lenght's
for idx, head := range data.Headers {
data.lenHeaders[idx] = visibleLen(head)
}
// determine max width per column
for _, entries := range data.rows {
currentWidth := 0
for idx, entry := range entries {
length := visibleLen(entry)
if data.lenHeaders[idx] < length {
if length > currentWidth+data.maxwidth {
data.lenHeaders[idx] = data.maxwidth - currentWidth
} else {
data.lenHeaders[idx] = length
}
}
currentWidth += data.lenHeaders[idx]
}
}
data.processed = true
} }
func isInt(num string) bool { func isInt(num string) bool {

View File

@@ -1,4 +1,4 @@
.PHONY: test clean cluster docs search up down delete .PHONY: test clean cluster docs search up down delete runenv
# docker stuff # docker stuff
up: up:
@@ -20,6 +20,8 @@ delete:
clean-docker: down delete clean-docker: down delete
runenv: up waitup cluster docs wait
# mosscap esctl stuff # mosscap esctl stuff
test: up waitup cluster docs wait search test: up waitup cluster docs wait search