Lands the executable half of the Phase 01-06 contracts: stdio MCP wiring with bundle policy enforced at call time, JWT/api_key upstream auth with single-flight renewal and 401 retry, strict argument decoding, allowlisted result projections with enum label twins, TV season expansion, settings read-modify-write under a revision lock, and a sanitized ToolError envelope that never forwards raw upstream bodies.
178 lines
5.0 KiB
Go
178 lines
5.0 KiB
Go
package tools
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"fmt"
|
|
|
|
"ombi-mcp/internal/ombi"
|
|
)
|
|
|
|
// permittedJobs is the closed enum of POST Job routes — names match
|
|
// the ledger exactly, including case-sensitive arrAvailability.
|
|
var permittedJobs = map[string]bool{
|
|
"update": true, "plexuserimporter": true, "plexwatchlist": true,
|
|
"embyuserimporter": true, "jellyfinuserimporter": true,
|
|
"plexcontentcacher": true, "clearmediaserverdata": true,
|
|
"plexrecentlyadded": true, "embycontentcacher": true,
|
|
"embyrecentlyadded": true, "jellyfincontentcacher": true,
|
|
"arrAvailability": true, "autodeleterequests": true,
|
|
"newsletter": true,
|
|
}
|
|
|
|
// write_job_run — trigger a permitted job or revalidate watchlist
|
|
// users. Accepted does not mean completed.
|
|
func handleJobRun(ctx context.Context, env *Env, raw json.RawMessage) *ToolResult {
|
|
o := newOp(ctx, env, "write_job_run", true)
|
|
var a JobRunArgs
|
|
if fail := o.args(raw, &a); fail != nil {
|
|
return fail
|
|
}
|
|
switch a.Action {
|
|
case "run":
|
|
if !permittedJobs[a.Job] {
|
|
return o.invalid("job", "unsupported job %q", a.Job)
|
|
}
|
|
raw, fail := o.call("POST", "/api/v1/Job/"+a.Job, nil, nil)
|
|
if fail != nil {
|
|
return fail
|
|
}
|
|
return o.jobResult(raw, a.Job)
|
|
case "revalidate_watchlist":
|
|
raw, fail := o.call("POST",
|
|
"/api/v1/Plex/WatchlistUsers/revalidate", nil, nil)
|
|
if fail != nil {
|
|
return fail
|
|
}
|
|
return o.jobResult(raw, "revalidate_watchlist")
|
|
default:
|
|
return o.invalid("action", "unsupported action %q", a.Action)
|
|
}
|
|
}
|
|
|
|
func (o *op) jobResult(raw []byte, name string) *ToolResult {
|
|
m := &Mutation{Kind: "mutation", Outcome: "accepted"}
|
|
if b, fail := o.decodeBool(raw); fail == nil {
|
|
m.UpstreamResult = &b
|
|
if !b {
|
|
return o.fail("UPSTREAM_REJECTED",
|
|
fmt.Sprintf("job %q was not accepted", name), false)
|
|
}
|
|
} else if obj, fail2 := o.decodeObject(raw); fail2 == nil {
|
|
if s := jstr(obj, "message"); s != "" {
|
|
m.Message = sanitizeText(s, maxSanitizedMsg)
|
|
}
|
|
}
|
|
return o.ok(m)
|
|
}
|
|
|
|
// write_notification_send — mass email to explicit resolved user
|
|
// IDs. Recipients are resolved privately; partial resolution failure
|
|
// is reported per item rather than silently dropped.
|
|
func handleNotificationSend(ctx context.Context, env *Env, raw json.RawMessage) *ToolResult {
|
|
o := newOp(ctx, env, "write_notification_send", true)
|
|
var a NotificationSendArgs
|
|
if fail := o.args(raw, &a); fail != nil {
|
|
return fail
|
|
}
|
|
if !nonempty(a.Subject) {
|
|
return o.invalid("subject", "subject must contain non-whitespace text")
|
|
}
|
|
if !nonempty(a.Body) {
|
|
return o.invalid("body", "body must contain non-whitespace text")
|
|
}
|
|
if len(a.UserIDs) == 0 {
|
|
return o.invalid("user_ids", "user_ids must be a nonempty list")
|
|
}
|
|
if len(a.UserIDs) > 100 {
|
|
return o.invalid("user_ids", "user_ids exceeds the 100-recipient cap")
|
|
}
|
|
seen := map[string]bool{}
|
|
users := []ombi.MassEmailUser{}
|
|
var itemResults []ItemResult
|
|
failed := 0
|
|
for _, id := range a.UserIDs {
|
|
if !nonempty(id) {
|
|
return o.invalid("user_ids", "user ids must be nonempty")
|
|
}
|
|
if seen[id] {
|
|
return o.invalid("user_ids", "duplicate user id %q", id)
|
|
}
|
|
seen[id] = true
|
|
raw, fail := o.call("GET", "/api/v1/Identity/User/"+seg(id), nil, nil)
|
|
if fail != nil {
|
|
failed++
|
|
itemResults = append(itemResults,
|
|
ItemResult{Identifier: id, Outcome: "rejected",
|
|
Message: "user lookup failed"})
|
|
continue
|
|
}
|
|
u, fail := o.decodeObject(raw)
|
|
if fail != nil {
|
|
failed++
|
|
itemResults = append(itemResults,
|
|
ItemResult{Identifier: id, Outcome: "rejected",
|
|
Message: "user record unreadable"})
|
|
continue
|
|
}
|
|
users = append(users, ombi.MassEmailUser{
|
|
ID: id,
|
|
UserName: jstr(u, "userName", "username"),
|
|
Email: jstr(u, "emailAddress", "email"),
|
|
})
|
|
itemResults = append(itemResults,
|
|
ItemResult{Identifier: id, Outcome: "accepted"})
|
|
}
|
|
if len(users) == 0 {
|
|
res := o.fail("PARTIAL_FAILURE",
|
|
"no recipients could be resolved", false)
|
|
return res
|
|
}
|
|
if failed > 0 {
|
|
o.warnf("%d of %d recipients could not be resolved", failed, len(a.UserIDs))
|
|
}
|
|
body := ombi.MassEmailModel{
|
|
Subject: a.Subject,
|
|
Body: a.Body,
|
|
Bcc: a.Bcc,
|
|
Users: users,
|
|
}
|
|
raw, fail := o.call("POST", "/api/v1/Notifications/massemail", nil, body)
|
|
if fail != nil {
|
|
return fail
|
|
}
|
|
m, fail := o.boolMutation(raw, nil)
|
|
if fail != nil {
|
|
return fail
|
|
}
|
|
if failed > 0 {
|
|
m.Outcome = "partial"
|
|
m.ItemResults = itemResults
|
|
}
|
|
n := len(users)
|
|
m.AffectedCount = &n
|
|
return o.ok(m)
|
|
}
|
|
|
|
// write_retry_remove — DELETE the queue ID only; never calls a
|
|
// nonexistent POST RequestRetry and does not touch the underlying
|
|
// request.
|
|
func handleRetryRemove(ctx context.Context, env *Env, raw json.RawMessage) *ToolResult {
|
|
o := newOp(ctx, env, "write_retry_remove", true)
|
|
var a RetryRemoveArgs
|
|
if fail := o.args(raw, &a); fail != nil {
|
|
return fail
|
|
}
|
|
if a.QueueID < 1 {
|
|
return o.invalid("queue_id", "queue_id must be a positive integer")
|
|
}
|
|
raw, fail := o.call("DELETE",
|
|
"/api/v1/RequestRetry/"+segInt(a.QueueID), nil, nil)
|
|
if fail != nil {
|
|
return fail
|
|
}
|
|
m := emptyMutation("completed", nil)
|
|
m.QueueID = &a.QueueID
|
|
return o.ok(m)
|
|
}
|