Skip to content

Commit 866c118

Browse files
Enhance IMAP handling in RunACTIONS function and add example configuration for email extraction
1 parent b701a0e commit 866c118

4 files changed

Lines changed: 73 additions & 4 deletions

File tree

.gitignore

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -40,4 +40,5 @@ logs.*.json
4040
*nyc*
4141
--public/
4242
examples/saas_model.md
43+
examples/downloads
4344
*.exe

examples/actions_email_ex.md

Lines changed: 56 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,56 @@
1+
# ACTIONS
2+
3+
```yaml metadata
4+
name: FileOperations
5+
description: "Transfer and organize generated reports"
6+
path: examples
7+
active: true
8+
```
9+
10+
## EMAILS
11+
```yaml
12+
name: emails
13+
description: Extract Emails
14+
type: IMAP
15+
params:
16+
protocol: IMAP
17+
host: imap.gmail.com
18+
port: 993
19+
username: "@SMTP_USERNAME"
20+
password: "@SMTP_PASSWORD"
21+
folder: INBOX
22+
download_att: true
23+
attachment_path: ./examples/downloads
24+
search:
25+
_from: supplier@example.com
26+
subject: C7 Intro
27+
since: 24h
28+
_before: 24h
29+
conn: "duckdb:"
30+
sqls:
31+
- ATTACH 'database/etl.db' AS DB (TYPE SQLITE)
32+
- create_emails
33+
- merge_into_emails
34+
- DETACH DB
35+
active: true
36+
```
37+
38+
```sql
39+
-- create_emails
40+
CREATE TABLE IF NOT EXISTS DB.emails (
41+
id BIGINT,
42+
subject VARCHAR,
43+
"from" VARCHAR,
44+
"to" VARCHAR,
45+
cc VARCHAR,
46+
bcc VARCHAR,
47+
date TIMESTAMP,
48+
body TEXT,
49+
attachments VARCHAR
50+
);
51+
```
52+
53+
```sql
54+
-- merge_into_emails
55+
INSERT INTO DB.emails BY NAME SELECT * FROM READ_JSON('<fname>')
56+
```

internal/etlx/imap.go

Lines changed: 13 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -16,8 +16,12 @@ import (
1616

1717
func returnAdresses(imapAdress []*imap.Address) string {
1818
adress := ""
19-
for _, adr := range imapAdress {
20-
adress += ";" + adr.Address()
19+
for i, adr := range imapAdress {
20+
glue := ""
21+
if i > 0 {
22+
glue = ";"
23+
}
24+
adress += glue + adr.Address()
2125
}
2226
return adress
2327
}
@@ -48,17 +52,20 @@ func (etlx *ETLX) ReadEmails(cfg map[string]any, item map[string]any, dateRef []
4852
// Connect IMAP
4953
c, err := client.DialTLS(fmt.Sprintf("%s:%s", host, port), nil)
5054
if err != nil {
55+
fmt.Println("client.DialTLS Err:", err)
5156
return nil, err
5257
}
5358
defer c.Logout()
5459
// Login
5560
err = c.Login(username, password)
5661
if err != nil {
62+
fmt.Println("c.Login Err:", err)
5763
return nil, err
5864
}
5965
// Select mailbox
6066
_, err = c.Select(folder, false)
6167
if err != nil {
68+
fmt.Println("c.Select(folder, false) Err:", err)
6269
return nil, err
6370
}
6471
// Build search
@@ -69,7 +76,8 @@ func (etlx *ETLX) ReadEmails(cfg map[string]any, item map[string]any, dateRef []
6976
case "from":
7077
criteria.Header.Add("From", fmt.Sprint(value))
7178
case "subject":
72-
criteria.Header.Add("Subject", fmt.Sprint(value))
79+
subj := etlx.SetQueryPlaceholders(fmt.Sprint(value), "", "", dateRef)
80+
criteria.Header.Add("Subject", subj)
7381
case "since":
7482
d, err := time.ParseDuration(fmt.Sprint(value))
7583
if err == nil {
@@ -85,8 +93,10 @@ func (etlx *ETLX) ReadEmails(cfg map[string]any, item map[string]any, dateRef []
8593
}
8694
ids, err := c.Search(criteria)
8795
if err != nil {
96+
fmt.Println("c.Search(criteria) Err:", err)
8897
return nil, err
8998
}
99+
// fmt.Println("Fine Till Search Criteria OK")
90100
results := []map[string]any{}
91101
for _, id := range ids {
92102
seq := new(imap.SeqSet)

internal/etlx/run_actions.go

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -468,7 +468,7 @@ func (etlx *ETLX) RunACTIONS(dateRef []time.Time, conf map[string]any, extraConf
468468
case "imap", "IMAP":
469469
_, okHost := params["host"].(string)
470470
_, okPass := params["password"].(string)
471-
_, okPort := params["port"].(string)
471+
_, okPort := params["port"]
472472
_, okUser := params["username"].(string)
473473
conn, okConn := params["conn"].(string)
474474
sqls, okSQLs := params["sqls"].([]any)
@@ -488,6 +488,7 @@ func (etlx *ETLX) RunACTIONS(dateRef []time.Time, conf map[string]any, extraConf
488488
_log2["msg"] = fmt.Sprintf("%s -> %s -> %s: missing required params port", key, itemKey, _type)
489489
valid = false
490490
}
491+
params["port"] = fmt.Sprintf("%v", params["port"])
491492
if !okUser {
492493
_log2["success"] = false
493494
_log2["msg"] = fmt.Sprintf("%s -> %s -> %s: missing required params user", key, itemKey, _type)
@@ -542,6 +543,7 @@ func (etlx *ETLX) RunACTIONS(dateRef []time.Time, conf map[string]any, extraConf
542543
}
543544
}
544545
}
546+
// fmt.Println("IMAP:", _log2["msg"])
545547
default:
546548
_log2["success"] = false
547549
_log2["msg"] = fmt.Sprintf("%s -> %s -> %s: Unsupported type", key, itemKey, _type)

0 commit comments

Comments
 (0)