FazBrowse GitHub Viewer
|
Trending
|
URL:
|
Home
Tools:
[Download Repo ZIP]
[View Raw Code]
[Original HTTPS Page]
odysseus/src/integrations.py at dev · odysseus-dev/odysseus · GitHub
odysseus-dev
/
odysseus
Public
Notifications
You must be signed in to change notification settings
Fork
1.3k
Star
90.9k
Code
Issues
1.1k
Pull requests
266
Discussions
Actions
Projects
Security and quality
0
Insights
Additional navigation options
Code
Issues
Pull requests
Discussions
Actions
Projects
Security and quality
Insights
Expand file tree
Breadcrumbs
odysseus
/
src
/
integrations.py
Copy path
More file actions
More file actions
Latest commit
History
History
History
810 lines (708 loc) · 33.5 KB
Breadcrumbs
odysseus
/
src
/
integrations.py
Copy path
File metadata and controls
810 lines (708 loc) · 33.5 KB
Raw
Copy raw file
Download raw file
Open symbols panel
Edit and raw actions
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
import
ipaddress
import
json
import
os
import
time
import
uuid
import
logging
import
re
from
typing
import
Dict
,
List
,
Optional
,
Any
from
urllib
.
parse
import
urljoin
,
urlparse
,
urlunparse
import
httpcore
import
httpx
from
fastapi
import
HTTPException
from
core
.
atomic_io
import
atomic_write_json
from
core
.
platform_compat
import
safe_chmod
from
src
.
secret_storage
import
decrypt
,
encrypt
,
is_encrypted
from
src
.
constants
import
DATA_DIR
,
INTEGRATIONS_FILE
,
SETTINGS_FILE
log
=
logging
.
getLogger
(
__name__
)
DATA_FILE
=
INTEGRATIONS_FILE
# ---------------------------------------------------------------------------
# Presets
# ---------------------------------------------------------------------------
INTEGRATION_PRESETS
:
Dict
[
str
,
Dict
[
str
,
Any
]]
=
{
"miniflux"
: {
"name"
:
"Miniflux"
,
"auth_type"
:
"header"
,
"auth_header"
:
"X-Auth-Token"
,
"description"
: (
"Miniflux RSS reader (v1 API). Key endpoints:
\n
"
" GET /v1/feeds — list all feeds
\n
"
" GET /v1/feeds/{id} — get feed details
\n
"
" POST /v1/feeds — create feed {
\"
feed_url
\"
:
\"
...
\"
,
\"
category_id
\"
: N}
\n
"
" PUT /v1/feeds/{id} — update feed
\n
"
" DELETE /v1/feeds/{id} — delete feed
\n
"
" GET /v1/feeds/{id}/entries — list entries for feed
\n
"
" GET /v1/entries — list all entries (params: status, limit, order, direction, category_id)
\n
"
" GET /v1/entries/{id} — get single entry
\n
"
" PUT /v1/entries — update entries {
\"
entry_ids
\"
: [...],
\"
status
\"
:
\"
read|unread
\"
}
\n
"
" GET /v1/categories — list categories
\n
"
" POST /v1/categories — create category {
\"
title
\"
:
\"
...
\"
}
\n
"
" GET /v1/feeds/{id}/icon — get feed icon
\n
"
" PUT /v1/entries/{id}/bookmark — toggle bookmark"
),
},
"gitea"
: {
"name"
:
"Gitea"
,
"auth_type"
:
"header"
,
"auth_header"
:
"Authorization"
,
"description"
: (
"Gitea git forge API (v1). Auth header value format: 'token YOUR_TOKEN'. Key endpoints:
\n
"
" GET /api/v1/repos/search — search repositories
\n
"
" GET /api/v1/repos/{owner}/{repo} — get repo details
\n
"
" GET /api/v1/repos/{owner}/{repo}/issues — list issues
\n
"
" POST /api/v1/repos/{owner}/{repo}/issues — create issue {
\"
title
\"
:
\"
...
\"
}
\n
"
" GET /api/v1/repos/{owner}/{repo}/pulls — list pull requests
\n
"
" GET /api/v1/repos/{owner}/{repo}/commits — list commits
\n
"
" GET /api/v1/user/repos — list your repos
\n
"
" GET /api/v1/orgs — list organizations
\n
"
" GET /api/v1/repos/{owner}/{repo}/contents/{filepath} — get file content"
),
},
"linkding"
: {
"name"
:
"Linkding"
,
"auth_type"
:
"header"
,
"auth_header"
:
"Authorization"
,
"description"
: (
"Linkding bookmark manager API. Auth header value format: 'Token YOUR_TOKEN'. Key endpoints:
\n
"
" GET /api/bookmarks/ — list bookmarks (params: q, limit, offset)
\n
"
" GET /api/bookmarks/{id}/ — get bookmark
\n
"
" POST /api/bookmarks/ — create bookmark {
\"
url
\"
:
\"
...
\"
,
\"
title
\"
:
\"
...
\"
,
\"
tag_names
\"
: [...]}
\n
"
" PUT /api/bookmarks/{id}/ — update bookmark
\n
"
" DELETE /api/bookmarks/{id}/ — delete bookmark
\n
"
" GET /api/bookmarks/archived/ — list archived bookmarks
\n
"
" GET /api/tags/ — list tags"
),
},
"homeassistant"
: {
"name"
:
"Home Assistant"
,
"auth_type"
:
"bearer"
,
"description"
: (
"Home Assistant smart home API. Key endpoints:
\n
"
" GET /api/ — API status check
\n
"
" GET /api/states — list all entity states
\n
"
" GET /api/states/{entity_id} — get state of entity
\n
"
" POST /api/states/{entity_id} — update entity state
\n
"
" POST /api/services/{domain}/{service} — call service (e.g. light/turn_on)
\n
"
" GET /api/history/period/{timestamp} — get state history
\n
"
" GET /api/logbook/{timestamp} — get logbook entries
\n
"
" POST /api/events/{event_type} — fire event
\n
"
" GET /api/config — get configuration"
),
},
"ntfy"
: {
"name"
:
"ntfy"
,
"auth_type"
:
"none"
,
"description"
: (
"ntfy push notification service. Key endpoints:
\n
"
" POST /{topic} — send notification. Body is the message text.
\n
"
" Headers: Title (notification title), Priority (1-5), Tags (comma-separated emoji tags)
\n
"
" POST / — send JSON notification {
\"
topic
\"
:
\"
...
\"
,
\"
message
\"
:
\"
...
\"
,
\"
title
\"
:
\"
...
\"
,
\"
priority
\"
: N}
\n
"
" GET /{topic}/json?poll=1 — poll for messages"
),
},
"discord_webhook"
: {
"name"
:
"Discord Webhook"
,
"auth_type"
:
"none"
,
"description"
: (
"Discord Incoming Webhook. Paste the full webhook URL (including the token) as the Base URL.
\n
"
"To get a URL: Discord server -> Server Settings -> Integrations -> Webhooks -> New Webhook -> Copy Webhook URL.
\n
"
"The secret is embedded in the URL — leave auth type as None.
\n
\n
"
"Use this integration as the target in Settings -> Reminders -> Webhook channel.
\n
"
"Payload template examples:
\n
"
" Simple: {
\"
content
\"
:
\"
{{title}}: {{message}}
\"
}
\n
"
" Embed: {
\"
embeds
\"
: [{
\"
title
\"
:
\"
{{title}}
\"
,
\"
description
\"
:
\"
{{message}}
\"
,
\"
color
\"
: 5793266}]}"
),
},
"vaultwarden"
: {
"name"
:
"Vaultwarden"
,
"auth_type"
:
"header"
,
"auth_header"
:
"Authorization"
,
"description"
: (
"Vaultwarden (Bitwarden-compatible) password manager API. Auth header value format: 'Bearer ACCESS_TOKEN'.
\n
"
"To get an access token: POST /identity/connect/token with grant_type=client_credentials&client_id=...&client_secret=...
\n
"
"Key endpoints:
\n
"
" GET /api/ciphers — list all vault items (logins, notes, cards, identities)
\n
"
" GET /api/ciphers/{id} — get a single vault item
\n
"
" POST /api/ciphers — create vault item {
\"
type
\"
: 1,
\"
name
\"
:
\"
...
\"
,
\"
login
\"
: {
\"
uri
\"
:
\"
...
\"
,
\"
username
\"
:
\"
...
\"
,
\"
password
\"
:
\"
...
\"
}}
\n
"
" PUT /api/ciphers/{id} — update vault item
\n
"
" DELETE /api/ciphers/{id} — delete vault item
\n
"
" GET /api/folders — list folders
\n
"
" POST /api/folders — create folder {
\"
name
\"
:
\"
...
\"
}
\n
"
" GET /api/collections — list collections (org vaults)
\n
"
" POST /api/ciphers/{id}/password-history — get password history
\n
"
" GET /api/sends — list Bitwarden Send items
\n
"
" POST /api/sends — create a Send (secure sharing)
\n
"
" Note: Vault data is end-to-end encrypted. The API returns encrypted fields
\n
"
" that must be decrypted client-side with the user's master key."
),
},
"freshrss"
: {
"name"
:
"FreshRSS"
,
"auth_type"
:
"header"
,
"auth_header"
:
"Authorization"
,
"description"
: (
"FreshRSS RSS reader (GReader API). Auth header value format: 'GoogleLogin auth=YOUR_TOKEN'. Key endpoints:
\n
"
" GET /api/greader.php/reader/api/0/subscription/list?output=json — list feeds
\n
"
" GET /api/greader.php/reader/api/0/stream/contents/feed/{feed_id}?output=json&n=20 — get entries
\n
"
" GET /api/greader.php/reader/api/0/tag/list?output=json — list tags/categories
\n
"
" POST /api/greader.php/reader/api/0/edit-tag — mark read/starred
\n
"
" GET /api/greader.php/reader/api/0/unread-count?output=json — unread counts"
),
},
}
# ---------------------------------------------------------------------------
# Storage
# ---------------------------------------------------------------------------
def
_ensure_data_dir
()
->
None
:
os
.
makedirs
(
os
.
path
.
dirname
(
DATA_FILE
),
exist_ok
=
True
)
def
_encrypt_integration_secrets
(
integrations
:
List
[
Dict
[
str
,
Any
]])
->
List
[
Dict
[
str
,
Any
]]:
"""Return storage-safe copies with API keys encrypted at rest."""
safe
:
List
[
Dict
[
str
,
Any
]]
=
[]
for
item
in
integrations
:
copy
=
dict
(
item
)
api_key
=
copy
.
get
(
"api_key"
,
""
)
if
api_key
:
copy
[
"api_key"
]
=
encrypt
(
str
(
api_key
))
safe
.
append
(
copy
)
return
safe
def
_decrypt_integration_secrets
(
integrations
:
List
[
Dict
[
str
,
Any
]])
->
List
[
Dict
[
str
,
Any
]]:
"""Return runtime copies with API keys decrypted for callers."""
decoded
:
List
[
Dict
[
str
,
Any
]]
=
[]
for
item
in
integrations
:
copy
=
dict
(
item
)
api_key
=
copy
.
get
(
"api_key"
,
""
)
if
api_key
:
copy
[
"api_key"
]
=
decrypt
(
str
(
api_key
))
decoded
.
append
(
copy
)
return
decoded
def
_has_plaintext_api_key
(
integrations
:
List
[
Dict
[
str
,
Any
]])
->
bool
:
return
any
(
bool
(
item
.
get
(
"api_key"
))
and
not
is_encrypted
(
str
(
item
.
get
(
"api_key"
)))
for
item
in
integrations
)
def
mask_integration_secret
(
integration
:
Dict
[
str
,
Any
])
->
Dict
[
str
,
Any
]:
"""Return a copy safe for API responses."""
safe
=
dict
(
integration
)
api_key
=
safe
.
get
(
"api_key"
,
""
)
if
api_key
:
safe
[
"api_key"
]
=
f"
{
str
(
api_key
)[:
4
]
}
****"
return
safe
def
_normalize_integration_base_url
(
base_url
:
Any
)
->
str
:
if
not
isinstance
(
base_url
,
str
)
or
not
base_url
.
strip
():
raise
ValueError
(
"Integration base URL is required"
)
cleaned
=
base_url
.
strip
().
rstrip
(
"/"
)
if
"?"
in
cleaned
or
"#"
in
cleaned
:
raise
ValueError
(
"Integration base URL must not include query or fragment"
)
parsed
=
urlparse
(
cleaned
)
if
parsed
.
scheme
.
lower
()
not
in
(
"http"
,
"https"
)
or
not
parsed
.
hostname
:
raise
ValueError
(
"Integration base URL must be an HTTP(S) URL"
)
return
urlunparse
(
parsed
.
_replace
(
scheme
=
parsed
.
scheme
.
lower
(),
query
=
""
,
fragment
=
""
)).
rstrip
(
"/"
)
def
_join_integration_url
(
base_url
:
str
,
path
:
str
)
->
str
:
base
=
base_url
.
rstrip
(
"/"
)
rel
=
path
.
lstrip
(
"/"
)
if
not
rel
:
# A bare "/" must resolve to the base URL itself, not base + "/".
# POST-to-base integrations (e.g. Discord webhooks) 404 on the
# trailing-slash variant of their URL.
return
base
return
urljoin
(
base
+
"/"
,
rel
)
def
load_integrations
()
->
List
[
Dict
[
str
,
Any
]]:
"""Load all integrations from disk with secrets decrypted for runtime use."""
if
not
os
.
path
.
exists
(
DATA_FILE
):
return
[]
try
:
with
open
(
DATA_FILE
,
"r"
,
encoding
=
"utf-8"
)
as
f
:
integrations
=
json
.
load
(
f
)
if
not
isinstance
(
integrations
,
list
):
log
.
error
(
"Invalid integrations file shape: expected a list"
)
return
[]
valid_integrations
=
[
item
for
item
in
integrations
if
isinstance
(
item
,
dict
)]
if
len
(
valid_integrations
)
!=
len
(
integrations
):
log
.
error
(
"Invalid integrations file rows: ignored non-object entries"
)
integrations
=
valid_integrations
if
_has_plaintext_api_key
(
integrations
):
save_integrations
(
_decrypt_integration_secrets
(
integrations
))
return
_decrypt_integration_secrets
(
integrations
)
except
(
json
.
JSONDecodeError
,
IOError
)
as
exc
:
log
.
error
(
"Failed to load integrations: %s"
,
exc
)
return
[]
def
save_integrations
(
integrations
:
List
[
Dict
[
str
,
Any
]])
->
None
:
"""Persist integrations list to disk with API keys encrypted at rest."""
_ensure_data_dir
()
atomic_write_json
(
DATA_FILE
,
_encrypt_integration_secrets
(
integrations
),
indent
=
2
)
safe_chmod
(
DATA_FILE
,
0o600
)
def
get_integration
(
integration_id
:
str
)
->
Optional
[
Dict
[
str
,
Any
]]:
"""Get a single integration by id."""
for
item
in
load_integrations
():
if
item
.
get
(
"id"
)
==
integration_id
:
return
item
return
None
def
add_integration
(
data
:
Dict
[
str
,
Any
])
->
Dict
[
str
,
Any
]:
"""Add a new integration. If 'preset' is given, merge preset defaults first."""
integration
:
Dict
[
str
,
Any
]
=
{}
preset_key
=
data
.
get
(
"preset"
)
if
preset_key
and
preset_key
in
INTEGRATION_PRESETS
:
integration
.
update
(
INTEGRATION_PRESETS
[
preset_key
])
integration
[
"preset"
]
=
preset_key
integration
.
update
(
data
)
integration
.
setdefault
(
"id"
,
uuid
.
uuid4
().
hex
[:
12
])
integration
.
setdefault
(
"enabled"
,
True
)
integration
.
setdefault
(
"auth_type"
,
"none"
)
integration
.
setdefault
(
"auth_header"
,
""
)
integration
.
setdefault
(
"auth_param"
,
""
)
integration
.
setdefault
(
"description"
,
""
)
integration
.
setdefault
(
"api_key"
,
""
)
integration
.
setdefault
(
"name"
,
""
)
integration
.
setdefault
(
"base_url"
,
""
)
if
not
isinstance
(
integration
.
get
(
"name"
),
str
)
or
not
integration
[
"name"
].
strip
():
raise
HTTPException
(
400
,
"Integration name is required"
)
try
:
integration
[
"base_url"
]
=
_normalize_integration_base_url
(
integration
.
get
(
"base_url"
))
except
ValueError
as
exc
:
raise
HTTPException
(
400
,
str
(
exc
))
from
exc
integrations
=
load_integrations
()
integrations
.
append
(
integration
)
save_integrations
(
integrations
)
return
integration
def
update_integration
(
integration_id
:
str
,
data
:
Dict
[
str
,
Any
])
->
Optional
[
Dict
[
str
,
Any
]]:
"""Update fields on an existing integration. Returns updated integration or None."""
data
=
dict
(
data
)
if
"name"
in
data
and
(
not
isinstance
(
data
[
"name"
],
str
)
or
not
data
[
"name"
].
strip
()):
raise
HTTPException
(
400
,
"Integration name is required"
)
if
"base_url"
in
data
:
try
:
data
[
"base_url"
]
=
_normalize_integration_base_url
(
data
[
"base_url"
])
except
ValueError
as
exc
:
raise
HTTPException
(
400
,
str
(
exc
))
from
exc
integrations
=
load_integrations
()
for
item
in
integrations
:
if
item
.
get
(
"id"
)
==
integration_id
:
data
.
pop
(
"id"
,
None
)
# prevent id change
item
.
update
(
data
)
save_integrations
(
integrations
)
return
item
return
None
def
delete_integration
(
integration_id
:
str
)
->
bool
:
"""Delete an integration by id. Returns True if found and deleted."""
integrations
=
load_integrations
()
original_len
=
len
(
integrations
)
integrations
=
[
i
for
i
in
integrations
if
i
.
get
(
"id"
)
!=
integration_id
]
if
len
(
integrations
)
<
original_len
:
save_integrations
(
integrations
)
return
True
return
False
# ---------------------------------------------------------------------------
# API execution
# ---------------------------------------------------------------------------
def
_strip_html_tags
(
html
:
str
)
->
str
:
"""Rough HTML tag stripping."""
text
=
re
.
sub
(
r"<[^>]+>"
,
""
,
html
)
text
=
re
.
sub
(
r"\s+"
,
" "
,
text
).
strip
()
return
text
def
_find_integration
(
identifier
:
str
)
->
Optional
[
Dict
[
str
,
Any
]]:
"""Find integration by id or name (case-insensitive)."""
integrations
=
load_integrations
()
# try id first
for
item
in
integrations
:
if
item
.
get
(
"id"
)
==
identifier
:
return
item
# try name
lower
=
identifier
.
lower
()
for
item
in
integrations
:
if
item
.
get
(
"name"
,
""
).
lower
()
==
lower
:
return
item
return
None
# httpcore raises its own exception hierarchy; map the ones a simple request can
# surface back to their httpx equivalents so the caller's `except httpx.*` blocks
# below behave exactly as they did with the default transport.
_HTTPCORE_TO_HTTPX_EXC
=
{
httpcore
.
ConnectError
:
httpx
.
ConnectError
,
httpcore
.
ConnectTimeout
:
httpx
.
ConnectTimeout
,
httpcore
.
NetworkError
:
httpx
.
NetworkError
,
httpcore
.
PoolTimeout
:
httpx
.
PoolTimeout
,
httpcore
.
ProtocolError
:
httpx
.
ProtocolError
,
httpcore
.
ReadError
:
httpx
.
ReadError
,
httpcore
.
ReadTimeout
:
httpx
.
ReadTimeout
,
httpcore
.
RemoteProtocolError
:
httpx
.
RemoteProtocolError
,
httpcore
.
TimeoutException
:
httpx
.
TimeoutException
,
httpcore
.
WriteError
:
httpx
.
WriteError
,
httpcore
.
WriteTimeout
:
httpx
.
WriteTimeout
,
}
class
_PinnedAsyncBackend
(
httpcore
.
AsyncNetworkBackend
):
"""Network backend that connects only to the pre-validated IPs, in order.
Every address here came out of the single SSRF resolution, so moving to the
next one after a connect failure is not re-resolution — it's ordinary
multi-address fallback restricted to the set the guard already approved.
httpcore takes TLS SNI and the ``Host`` header from the request URL rather
than the connect host, so pinning the socket destination leaves certificate
validation and vhost routing pointed at the original hostname.
"""
def
__init__
(
self
,
ips
:
List
[
ipaddress
.
_BaseAddress
]):
self
.
_ips
=
[
str
(
ip
)
for
ip
in
ips
]
self
.
_real
=
httpcore
.
AnyIOBackend
()
async
def
connect_tcp
(
self
,
host
,
port
,
timeout
=
None
,
local_address
=
None
,
socket_options
=
None
):
# One shared connect budget: each attempt gets the time left until the
# original deadline, so N dead addresses can't stretch the connect
# phase to N * timeout.
deadline
=
None
if
timeout
is
None
else
time
.
monotonic
()
+
timeout
last_exc
:
Optional
[
Exception
]
=
None
for
ip
in
self
.
_ips
:
remaining
=
None
if
deadline
is
None
else
max
(
0.0
,
deadline
-
time
.
monotonic
())
try
:
return
await
self
.
_real
.
connect_tcp
(
ip
,
port
,
remaining
,
local_address
,
socket_options
)
except
(
httpcore
.
ConnectError
,
httpcore
.
ConnectTimeout
)
as
exc
:
last_exc
=
exc
if
deadline
is
not
None
and
time
.
monotonic
()
>=
deadline
:
break
raise
last_exc
async
def
connect_unix_socket
(
self
,
path
,
timeout
=
None
,
socket_options
=
None
):
return
await
self
.
_real
.
connect_unix_socket
(
path
,
timeout
,
socket_options
)
async
def
sleep
(
self
,
seconds
:
float
)
->
None
:
return
await
self
.
_real
.
sleep
(
seconds
)
class
_PinnedAsyncTransport
(
httpx
.
AsyncBaseTransport
):
"""httpx transport that pins the TCP connect to the pre-resolved IP(s).
Kept local, mirroring the per-module pinned transports web fetch and
webhook delivery already carry, rather than coupling api_call to the
webhook subsystem. The request URL passes through unchanged, so SNI and the
``Host`` header stay the original hostname; only the socket destination is
pinned, which is what closes the rebinding window.
"""
def
__init__
(
self
,
ips
:
List
[
ipaddress
.
_BaseAddress
]):
self
.
_pinned_ips
=
list
(
ips
)
self
.
_pool
=
httpcore
.
AsyncConnectionPool
(
# Reuse the CA trust the default httpx client would build (certifi
# plus SSL_CERT_FILE / SSL_CERT_DIR when trust_env is set) so
# swapping in this transport doesn't quietly change which chains
# verify. ssl.create_default_context() would use system roots.
ssl_context
=
httpx
.
create_ssl_context
(),
http1
=
True
,
http2
=
False
,
network_backend
=
_PinnedAsyncBackend
(
ips
),
)
async
def
handle_async_request
(
self
,
request
:
httpx
.
Request
)
->
httpx
.
Response
:
core_req
=
httpcore
.
Request
(
method
=
request
.
method
,
url
=
httpcore
.
URL
(
scheme
=
request
.
url
.
raw_scheme
,
host
=
request
.
url
.
raw_host
,
port
=
request
.
url
.
port
,
target
=
request
.
url
.
raw_path
,
),
headers
=
request
.
headers
.
raw
,
content
=
request
.
stream
,
extensions
=
request
.
extensions
,
)
try
:
core_resp
=
await
self
.
_pool
.
handle_async_request
(
core_req
)
content
=
b""
.
join
([
chunk
async
for
chunk
in
core_resp
.
aiter_stream
()])
await
core_resp
.
aclose
()
except
Exception
as
exc
:
mapped
=
_HTTPCORE_TO_HTTPX_EXC
.
get
(
type
(
exc
))
if
mapped
is
not
None
:
raise
mapped
(
str
(
exc
))
from
exc
raise
return
httpx
.
Response
(
status_code
=
core_resp
.
status
,
headers
=
core_resp
.
headers
,
content
=
content
,
extensions
=
core_resp
.
extensions
,
)
async
def
aclose
(
self
)
->
None
:
await
self
.
_pool
.
aclose
()
def
_validated_ips
(
raw_ips
:
List
[
str
])
->
List
[
ipaddress
.
_BaseAddress
]:
"""Return every entry that parses as an IP address, de-duplicated, order
preserved.
check_outbound_url only reports ok when *all* of these classify as safe, so
the whole list is guard-approved and any of them is a legitimate connect
target. Skipping unparseable entries mirrors how the guard walks the same
resolver output.
De-duplication matters because the resolver is getaddrinfo(host, None) with
no socktype filter, so glibc reports the same address once per socktype
(SOCK_STREAM/SOCK_DGRAM/SOCK_RAW) — a single-homed host comes back three
times. Without this, the connect fallback would spend the shared deadline
retrying one dead address instead of moving on to a genuinely different one.
"""
ips
:
List
[
ipaddress
.
_BaseAddress
]
=
[]
seen
=
set
()
for
raw
in
raw_ips
:
if
not
isinstance
(
raw
,
str
):
continue
try
:
ip
=
ipaddress
.
ip_address
(
raw
.
split
(
"%"
)[
0
])
# strip IPv6 zone id
except
ValueError
:
continue
if
ip
in
seen
:
continue
seen
.
add
(
ip
)
ips
.
append
(
ip
)
return
ips
async
def
execute_api_call
(
integration_id
:
str
,
method
:
str
,
path
:
str
,
params
:
Optional
[
Dict
[
str
,
Any
]]
=
None
,
body
:
Optional
[
Any
]
=
None
,
extra_headers
:
Optional
[
Dict
[
str
,
str
]]
=
None
,
)
->
Dict
[
str
,
Any
]:
"""Execute an HTTP request against a registered integration."""
integration
=
_find_integration
(
integration_id
)
if
not
integration
:
return
{
"error"
:
f"Integration not found:
{
integration_id
}
"
,
"exit_code"
:
1
}
if
not
integration
.
get
(
"enabled"
,
True
):
return
{
"error"
:
f"Integration '
{
integration
.
get
(
'name'
)
}
' is disabled"
,
"exit_code"
:
1
}
try
:
base_url
=
_normalize_integration_base_url
(
integration
.
get
(
"base_url"
,
""
))
except
ValueError
as
exc
:
return
{
"error"
:
str
(
exc
),
"exit_code"
:
1
}
# Strip common API path suffixes users might accidentally include
# (e.g. "http://host/v1/" → "http://host"). The integration's preset
# endpoints include the full path, so the base should be bare.
preset
=
(
integration
.
get
(
"preset"
)
or
integration
.
get
(
"name"
,
""
)).
lower
()
strip_suffixes
=
{
"miniflux"
: [
"/v1"
],
"gitea"
: [
"/api/v1"
,
"/api"
],
"linkding"
: [
"/api"
],
"homeassistant"
: [
"/api"
],
}
for
suf
in
strip_suffixes
.
get
(
preset
, []):
if
base_url
.
endswith
(
suf
):
base_url
=
base_url
[:
-
len
(
suf
)]
break
# Validate path
if
not
path
.
startswith
(
"/"
):
return
{
"error"
:
"Path must start with /"
,
"exit_code"
:
1
}
if
re
.
search
(
r"^https?://"
,
path
)
or
"://"
in
path
:
return
{
"error"
:
"Path must not contain a protocol scheme"
,
"exit_code"
:
1
}
if
"#"
in
path
:
return
{
"error"
:
"Path must not contain a fragment"
,
"exit_code"
:
1
}
url
=
_join_integration_url
(
base_url
,
path
)
# SSRF guard — same check used by the gallery endpoint, embeddings,
# CardDAV, and the reminder webhook sender. Link-local / metadata
# addresses (169.254.x.x — the cloud credential-exfil vector) are always
# rejected; INTEGRATION_API_BLOCK_PRIVATE_IPS=true also blocks RFC-1918 /
# loopback for locked-down deployments. Private stays allowed by default
# because LAN integrations (Home Assistant, Miniflux, ntfy) are the
# primary use case.
from
src
.
url_safety
import
check_outbound_url
,
_default_resolver
block_private
=
os
.
getenv
(
"INTEGRATION_API_BLOCK_PRIVATE_IPS"
,
"false"
).
lower
()
==
"true"
# Resolve the host exactly once and remember the IPs the guard validated so
# the request below can be pinned to them. check_outbound_url only reports
# (ok, reason); a plain httpx client re-resolves the host at connect time,
# which reopens a DNS-rebinding TOCTOU — a base_url host that answers with a
# public IP for the guard and then flips to 169.254.169.254 for the connect
# would reach cloud metadata with the integration's auth headers attached.
resolved_ips
:
List
[
str
]
=
[]
def
_recording_resolver
(
host
:
str
)
->
List
[
str
]:
ips
=
_default_resolver
(
host
)
resolved_ips
[:]
=
ips
return
ips
ok
,
reason
=
check_outbound_url
(
url
,
block_private
=
block_private
,
resolver
=
_recording_resolver
)
if
not
ok
:
return
{
"error"
:
f"URL rejected:
{
reason
}
"
,
"exit_code"
:
1
}
pinned_ips
=
_validated_ips
(
resolved_ips
)
if
not
pinned_ips
:
return
{
"error"
:
"URL rejected: host did not resolve to a usable address"
,
"exit_code"
:
1
}
method
=
method
.
upper
()
# Build headers
headers
:
Dict
[
str
,
str
]
=
{}
if
extra_headers
:
headers
.
update
(
extra_headers
)
api_key
=
integration
.
get
(
"api_key"
,
""
)
auth_type
=
integration
.
get
(
"auth_type"
,
"none"
)
if
auth_type
==
"header"
and
api_key
:
# Fall back based on preset/name when auth_header is unset or empty
header_name
=
integration
.
get
(
"auth_header"
)
or
""
if
not
header_name
:
preset
=
(
integration
.
get
(
"preset"
)
or
integration
.
get
(
"name"
,
""
)).
lower
()
header_defaults
=
{
"miniflux"
:
"X-Auth-Token"
,
"linkding"
:
"Authorization"
,
"gitea"
:
"Authorization"
,
}
header_name
=
header_defaults
.
get
(
preset
,
"Authorization"
)
headers
[
header_name
]
=
api_key
elif
auth_type
==
"bearer"
and
api_key
:
headers
[
"Authorization"
]
=
f"Bearer
{
api_key
}
"
elif
auth_type
==
"query"
and
api_key
:
if
params
is
None
:
params
=
{}
param_name
=
integration
.
get
(
"auth_param"
,
"api_key"
)
params
[
param_name
]
=
api_key
# auth_type == "basic" — expects api_key as "user:password"
auth
=
None
if
auth_type
==
"basic"
and
api_key
:
parts
=
api_key
.
split
(
":"
,
1
)
if
len
(
parts
)
==
2
:
auth
=
httpx
.
BasicAuth
(
parts
[
0
],
parts
[
1
])
try
:
async
with
httpx
.
AsyncClient
(
timeout
=
30.0
,
transport
=
_PinnedAsyncTransport
(
pinned_ips
)
)
as
client
:
response
=
await
client
.
request
(
method
,
url
,
params
=
params
,
json
=
body
if
body
is
not
None
else
None
,
headers
=
headers
,
auth
=
auth
,
)
content_type
=
response
.
headers
.
get
(
"content-type"
,
""
)
status
=
response
.
status_code
# Format response body
if
"application/json"
in
content_type
:
try
:
data
=
response
.
json
()
full
=
json
.
dumps
(
data
,
indent
=
2
,
ensure_ascii
=
False
)
if
len
(
full
)
>
12000
:
if
isinstance
(
data
,
list
):
# Binary-search for the largest prefix such that the
# final array (prefix + sentinel) fits within the limit.
# Pre-compute the sentinel so we know its serialized size.
sentinel_placeholder
=
{
"_truncated"
:
True
,
"total_items"
:
len
(
data
),
"shown_items"
:
0
,
}
# Overhead: the sentinel appears as an extra array element.
# Add a conservative padding for the separating comma,
# newline, and indentation characters (~6 chars).
sentinel_overhead
=
len
(
json
.
dumps
(
sentinel_placeholder
,
indent
=
2
,
ensure_ascii
=
False
)
)
+
6
budget
=
12000
-
sentinel_overhead
lo
,
hi
=
0
,
len
(
data
)
while
lo
<
hi
:
mid
=
(
lo
+
hi
+
1
)
//
2
candidate
=
json
.
dumps
(
data
[:
mid
],
indent
=
2
,
ensure_ascii
=
False
)
if
len
(
candidate
)
<
budget
:
lo
=
mid
else
:
hi
=
mid
-
1
sentinel
=
{
"_truncated"
:
True
,
"total_items"
:
len
(
data
),
"shown_items"
:
lo
,
}
formatted
=
json
.
dumps
(
data
[:
lo
]
+
[
sentinel
],
indent
=
2
,
ensure_ascii
=
False
)
elif
isinstance
(
data
,
dict
):
# Truncate dict entries until the result fits, then add
# the _truncated marker. Walk keys in insertion order.
DICT_LIMIT
=
12000
kept
:
dict
=
{}
for
k
,
v
in
data
.
items
():
candidate
=
json
.
dumps
(
{
**
kept
,
k
:
v
,
"_truncated"
:
True
},
indent
=
2
,
ensure_ascii
=
False
,
)
if
len
(
candidate
)
<=
DICT_LIMIT
:
kept
[
k
]
=
v
else
:
break
formatted
=
json
.
dumps
(
{
**
kept
,
"_truncated"
:
True
},
indent
=
2
,
ensure_ascii
=
False
)
else
:
total
=
len
(
full
)
formatted
=
full
[:
12000
]
+
f"
\n
... (truncated,
{
total
}
chars total)"
else
:
formatted
=
full
except
(
json
.
JSONDecodeError
,
ValueError
):
formatted
=
response
.
text
if
len
(
formatted
)
>
12000
:
total
=
len
(
formatted
)
formatted
=
formatted
[:
12000
]
+
f"
\n
... (truncated,
{
total
}
chars total)"
elif
"text/html"
in
content_type
:
formatted
=
_strip_html_tags
(
response
.
text
)
if
len
(
formatted
)
>
12000
:
total
=
len
(
formatted
)
formatted
=
formatted
[:
12000
]
+
f"
\n
... (truncated,
{
total
}
chars total)"
else
:
formatted
=
response
.
text
if
len
(
formatted
)
>
12000
:
total
=
len
(
formatted
)
formatted
=
formatted
[:
12000
]
+
f"
\n
... (truncated,
{
total
}
chars total)"
output
=
f"HTTP
{
status
}
\n
{
formatted
}
"
if
status
>=
400
:
return
{
"error"
:
output
,
"exit_code"
:
1
,
# The error string includes the remote response body. Preserve
# it for diagnostics, but make its provenance explicit so the
# agent gate does not treat HTTP failure as content-free.
"untrusted_content"
:
True
,
}
return
{
"output"
:
output
,
"exit_code"
:
0
}
except
httpx
.
TimeoutException
:
return
{
"error"
:
f"Request to
{
integration
.
get
(
'name'
)
}
timed out"
,
"exit_code"
:
1
}
except
httpx
.
RequestError
as
exc
:
return
{
"error"
:
f"Request failed:
{
exc
}
"
,
"exit_code"
:
1
}
except
Exception
as
exc
:
log
.
exception
(
"Unexpected error in execute_api_call"
)
return
{
"error"
:
f"Unexpected error:
{
exc
}
"
,
"exit_code"
:
1
}
# ---------------------------------------------------------------------------
# System prompt helper
# ---------------------------------------------------------------------------
def
get_integrations_prompt
()
->
str
:
"""Return a string describing all enabled integrations for system prompt injection.
Returns empty string if no integrations are enabled.
"""
integrations
=
load_integrations
()
enabled
=
[
i
for
i
in
integrations
if
i
.
get
(
"enabled"
,
True
)]
if
not
enabled
:
return
""
lines
=
[
"You have access to the following API integrations via the api_call tool:
\n
"
]
for
integ
in
enabled
:
name
=
integ
.
get
(
"name"
,
integ
.
get
(
"id"
,
"unknown"
))
lines
.
append
(
f"##
{
name
}
(id:
{
integ
[
'id'
]
}
)"
)
desc
=
integ
.
get
(
"description"
,
""
)
if
desc
:
lines
.
append
(
desc
)
lines
.
append
(
""
)
return
"
\n
"
.
join
(
lines
)
# ---------------------------------------------------------------------------
# Migration
# ---------------------------------------------------------------------------
def
migrate_from_settings
()
->
None
:
"""If data/settings.json has miniflux_url and miniflux_api_key, create a
Miniflux integration and clear those keys from settings."""
settings_path
=
SETTINGS_FILE
if
not
os
.
path
.
exists
(
settings_path
):
return
try
:
with
open
(
settings_path
,
"r"
,
encoding
=
"utf-8"
)
as
f
:
settings
=
json
.
load
(
f
)
except
(
json
.
JSONDecodeError
,
IOError
):
return
miniflux_url
=
settings
.
get
(
"miniflux_url"
,
""
)
miniflux_key
=
settings
.
get
(
"miniflux_api_key"
,
""
)
if
not
miniflux_url
or
not
miniflux_key
:
return
# Check if a miniflux integration already exists
existing
=
load_integrations
()
for
item
in
existing
:
if
item
.
get
(
"preset"
)
==
"miniflux"
:
log
.
info
(
"Miniflux integration already exists, skipping migration"
)
return
add_integration
({
"preset"
:
"miniflux"
,
"base_url"
:
miniflux_url
.
rstrip
(
"/"
),
"api_key"
:
miniflux_key
,
})
# Clear migrated keys
settings
.
pop
(
"miniflux_url"
,
None
)
settings
.
pop
(
"miniflux_api_key"
,
None
)
with
open
(
settings_path
,
"w"
,
encoding
=
"utf-8"
)
as
f
:
json
.
dump
(
settings
,
f
,
indent
=
2
)
log
.
info
(
"Migrated Miniflux integration from settings.json"
)
Back
|
FazBrowse Home
|
New Git URL