"""Jobber connector.
Syncs clients, quotes (estimates), jobs, and invoices from Jobber's
GraphQL API using OAuth 2.0 authorization code flow.
Data mapping:
- Jobber Client → JobberClient
- Jobber Quote → JobberQuote → also creates Estimate record
- Jobber Job → JobberJob
- Jobber Invoice → JobberInvoice
GraphQL API docs: https://developer.getjobber.com/api/
"""
from __future__ import annotations
import logging
import time
from datetime import datetime, timezone, timedelta
from typing import Any, Dict, List, Optional
import requests
from . import BaseConnector, OAuthConnector, _REGISTRY, register_connector
logger = logging.getLogger(__name__)
# -- Registration metadata ---------------------------------------------------
register_connector(
"jobber",
{
"service": "jobber",
"name": "Jobber",
"category": "crm",
"description": "Field service management — clients, quotes, jobs, invoices via GraphQL.",
"auth_type": "oauth2",
"auth_fields": [
"client_id",
"client_secret",
"access_token",
"refresh_token",
],
"capabilities": ["clients", "quotes", "jobs", "invoices", "estimates"],
"rate_limit": "5000 req/hr",
"docs_url": "https://developer.getjobber.com/api/",
},
)
# -- Constants ---------------------------------------------------------------
_TOKEN_URL = "https://api.jobber.com/oauth/token"
_GRAPHQL_URL = "https://api.getjobber.com/api/graphql"
_GRAPHQL_VERSION = "2023-11-15"
_AUTHORIZE_URL = "https://api.jobber.com/oauth/authorize"
_SCOPES = [
"read", # Read access to all data
"write", # Write access (optional)
]
# -- GraphQL queries ---------------------------------------------------------
CLIENTS_QUERY = """
query clients($after: String, $limit: Int) {
clients(after: $after, limit: $limit) {
pageInfo {
endCursor
hasNextPage
}
edges {
node {
id
firstName
lastName
companyName
defaultAddress {
line1
line2
city
province
country
postalCode
}
contacts {
id
email
mobilePhone
}
timezones {
id
}
languageCode
}
}
}
}
"""
QUOTES_QUERY = """
query quotes($after: String, $limit: Int) {
quotes(after: $after, limit: $limit) {
pageInfo {
endCursor
hasNextPage
}
edges {
node {
id
name
status
quoteValue
currencyCode
client {
id
}
sentAt
acceptedAt
rejectedAt
expiresAt
lineItems {
id
name
description
value
quantity
unit
}
tax {
total
}
discount {
total
}
}
}
}
}
"""
JOBS_QUERY = """
query jobs($after: String, $limit: Int) {
jobs(after: $after, limit: $limit) {
pageInfo {
endCursor
hasNextPage
}
edges {
node {
id
name
description
status
type
jobValue {
estimated
actual
}
jobCost {
estimated
actual
}
scheduledStart
scheduledEnd
actualStart
actualEnd
client {
id
}
address {
line1
line2
city
province
country
postalCode
}
}
}
}
}
"""
INVOICES_QUERY = """
query invoices($after: String, $limit: Int) {
invoices(after: $after, limit: $limit) {
pageInfo {
endCursor
hasNextPage
}
edges {
node {
id
number
status
total
currencyCode
taxTotal
paidTotal
client {
id
}
job {
id
}
dueDate
createdAt
paidAt
lineItems {
id
name
quantity
unit
value
}
}
}
}
}
"""
# -- JobberConnector ---------------------------------------------------------
class JobberConnector(OAuthConnector):
"""Jobber integration — pulls clients, quotes, jobs, invoices via GraphQL."""
_SERVICE = "jobber"
OAUTH_AUTHORIZE_URL = _AUTHORIZE_URL
OAUTH_TOKEN_URL = _TOKEN_URL
OAUTH_SCOPES = _SCOPES
# Token cache
_access_token: Optional[str] = None
_refresh_token: Optional[str] = None
_token_expires_at: Optional[datetime] = None
PAGE_LIMIT = 50
def __init__(self, *, company_id: str, config: Dict[str, Any], connector_id: Optional[str] = None) -> None:
super().__init__(company_id=company_id, config=config, connector_id=connector_id)
# Pre-load token from config if available
if self.config.get("access_token") and self.config.get("refresh_token"):
self._access_token = self.config["access_token"]
self._refresh_token = self.config["refresh_token"]
self._token_expires_at = datetime.now(timezone.utc) + timedelta(hours=1)
# -- Token management ----------------------------------------------------
def _get_client_id(self) -> str:
"""Read OAuth client ID from config or Flask config."""
# Prefer tenant config over Flask global config
if self.config.get("client_id"):
return self.config["client_id"]
try:
return super()._get_client_id()
except Exception:
return ""
def _get_client_secret(self) -> str:
"""Read OAuth client secret from config or Flask config."""
if self.config.get("client_secret"):
return self.config["client_secret"]
try:
return super()._get_client_secret()
except Exception:
return ""
def _get_redirect_uri(self) -> str:
"""Read redirect URI from config or Flask config."""
if self.config.get("redirect_uri"):
return self.config["redirect_uri"]
try:
return super()._get_redirect_uri()
except Exception:
return ""
def _get_access_token(self) -> str:
"""Return a valid access token, refreshing if needed."""
now = datetime.now(timezone.utc)
if (
self._access_token
and self._token_expires_at
and now < self._token_expires_at - timedelta(seconds=60)
):
return self._access_token
# Try config-stored token first
config_token = self.config.get("access_token", "")
config_refresh = self.config.get("refresh_token", "")
if config_token and config_refresh:
self._access_token = config_token
self._refresh_token = config_refresh
# Assume 1hr expiry if not tracked
self._token_expires_at = now + timedelta(hours=1)
return self._access_token
raise ValueError("No access token available. Complete OAuth flow first.")
def _refresh_token_impl(self) -> None:
"""Refresh access token using stored refresh_token."""
refresh = self.config.get("refresh_token") or self._refresh_token
if not refresh:
raise ValueError("No refresh token available")
client_id = self._get_client_id()
client_secret = self._get_client_secret()
resp = requests.post(
_TOKEN_URL,
data={
"grant_type": "refresh_token",
"refresh_token": refresh,
"client_id": client_id,
"client_secret": client_secret,
},
timeout=30,
)
resp.raise_for_status()
data = resp.json()
self._access_token = data.get("access_token")
self._refresh_token = data.get("refresh_token", refresh)
expires_in = int(data.get("expires_in", 3600))
self._token_expires_at = datetime.now(timezone.utc) + timedelta(seconds=expires_in)
# Persist in config
self.config["access_token"] = self._access_token
self.config["refresh_token"] = self._refresh_token
logger.info("Jobber token refreshed, expires in %ds", expires_in)
# -- GraphQL helpers -----------------------------------------------------
def _graphql_request(self, query: str, variables: Optional[Dict] = None) -> Dict:
"""Execute a GraphQL query with pagination support."""
token = self._get_access_token()
resp = requests.post(
_GRAPHQL_URL,
json={"query": query, "variables": variables or {}},
headers={
"Authorization": f"Bearer {token}",
"Content-Type": "application/json",
"X-JOBBER-GRAPHQL-VERSION": _GRAPHQL_VERSION,
},
timeout=60,
)
resp.raise_for_status()
data = resp.json()
if "errors" in data:
errors = data["errors"]
if isinstance(errors, list) and errors:
error_msg = errors[0].get("message", str(errors))
raise ValueError(f"Jobber GraphQL error: {error_msg}")
return data.get("data", {})
def _paginated_query(self, query: str, entity_key: str) -> List[Dict]:
"""Handle cursor-based pagination for a GraphQL query."""
all_items: List[Dict] = []
after = None
while True:
variables = {"limit": self.PAGE_LIMIT}
if after:
variables["after"] = after
data = self._graphql_request(query, variables)
collection = data.get(entity_key, {})
edges = collection.get("edges", [])
page_info = collection.get("pageInfo", {})
for edge in edges:
all_items.append(edge.get("node", {}))
if not page_info.get("hasNextPage", False):
break
after = page_info.get("endCursor")
return all_items
# -- connect / disconnect ------------------------------------------------
def connect(self) -> Dict[str, Any]:
"""Validate credentials by fetching a client."""
self._log(event_type="connect_attempt", status="pending")
start = time.monotonic()
required = ["client_id", "client_secret", "access_token"]
missing = [f for f in required if not self.config.get(f)]
if missing:
error = f"Missing required fields: {', '.join(missing)}"
self._log(event_type="connect_error", status="error", error_message=error)
return {"status": "error", "error": error}
try:
clients = self._paginated_query(CLIENTS_QUERY, "clients")
duration_ms = int((time.monotonic() - start) * 1000)
self._connected = True
self._log(
event_type="connected",
status="success",
record_count=len(clients),
duration_ms=duration_ms,
)
return {
"status": "connected",
"client_count": len(clients),
"duration_ms": duration_ms,
}
except Exception as exc:
duration_ms = int((time.monotonic() - start) * 1000)
self._log(
event_type="connect_error",
status="error",
error_message=str(exc),
duration_ms=duration_ms,
)
return {"status": "error", "error": str(exc)}
def disconnect(self) -> Dict[str, Any]:
"""Revoke tokens and clear local state."""
self._connected = False
self._access_token = None
self._refresh_token = None
self._token_expires_at = None
# Revoke token on server
if self.config.get("access_token"):
try:
client_id = self._get_client_id()
client_secret = self._get_client_secret()
requests.post(
"https://api.jobber.com/oauth/revoke",
data={
"token": self.config.get("access_token"),
"client_id": client_id,
"client_secret": client_secret,
},
timeout=15,
)
except Exception:
logger.warning("Jobber token revocation failed (non-critical)")
return {"status": "disconnected"}
# -- sync ----------------------------------------------------------------
def sync(self) -> Dict[str, Any]:
"""Sync clients, quotes, jobs, and invoices from Jobber."""
self._log(event_type="sync_start", status="pending")
start = time.monotonic()
results = {
"clients": 0,
"quotes": 0,
"jobs": 0,
"invoices": 0,
"estimates_created": 0,
}
try:
self._sync_clients()
self._sync_quotes()
self._sync_jobs()
self._sync_invoices()
# Update counts from DB
from app.models import (
db, JobberClient, JobberQuote, JobberJob, JobberInvoice,
)
results["clients"] = JobberClient.query.filter_by(
company_id=self.company_id
).count()
results["quotes"] = JobberQuote.query.filter_by(
company_id=self.company_id
).count()
results["jobs"] = JobberJob.query.filter_by(
company_id=self.company_id
).count()
results["invoices"] = JobberInvoice.query.filter_by(
company_id=self.company_id
).count()
duration_ms = int((time.monotonic() - start) * 1000)
self._log(
event_type="sync_complete",
status="success",
record_count=sum(results.values()),
duration_ms=duration_ms,
details=results,
)
return {
"status": "success",
"records": results,
"duration_ms": duration_ms,
}
except Exception as exc:
duration_ms = int((time.monotonic() - start) * 1000)
self._log(
event_type="sync_error",
status="error",
error_message=str(exc),
duration_ms=duration_ms,
)
return {"status": "error", "error": str(exc)}
def _sync_clients(self):
"""Sync clients from Jobber."""
raw_clients = self._paginated_query(CLIENTS_QUERY, "clients")
logger.info("Jobber: fetched %d clients", len(raw_clients))
from app.models import db, JobberClient
now = datetime.now(timezone.utc)
for client in raw_clients:
external_id = client.get("id", "")
if not external_id:
continue
# Parse contacts
contacts = client.get("contacts", [])
primary_contact = contacts[0] if contacts else {}
# Parse address
address = client.get("defaultAddress", {}) or {}
existing = JobberClient.query.filter_by(
company_id=self.company_id,
external_id=external_id,
).first()
data = {
"external_id": external_id,
"first_name": client.get("firstName", "") or "",
"last_name": client.get("lastName", "") or "",
"company_name": client.get("companyName", "") or "",
"email": primary_contact.get("email", "") or "",
"phone": primary_contact.get("mobilePhone", "") or "",
"address_line1": address.get("line1", "") or "",
"address_line2": address.get("line2", "") or "",
"city": address.get("city", "") or "",
"state": address.get("province", "") or "",
"postal_code": address.get("postalCode", "") or "",
"country": address.get("country", "US") or "US",
"jobber_timezone": (client.get("timezones") or [{}])[0].get("id", "") if client.get("timezones") else "",
"jobber_language": client.get("languageCode", "en") or "en",
"status": "active",
"synced_at": now,
}
if existing:
existing.__dict__.update(data)
else:
data["company_id"] = self.company_id
new = JobberClient(**data)
db.session.add(new)
db.session.commit()
def _sync_quotes(self):
"""Sync quotes from Jobber and create linked Estimates."""
raw_quotes = self._paginated_query(QUOTES_QUERY, "quotes")
logger.info("Jobber: fetched %d quotes", len(raw_quotes))
from app.models import db, JobberQuote, Estimate
now = datetime.now(timezone.utc)
for quote in raw_quotes:
external_id = quote.get("id", "")
if not external_id:
continue
client_data = quote.get("client", {}) or {}
client_external = client_data.get("id", "")
# Parse line items
line_items = quote.get("lineItems", []) or []
line_items_json = []
for item in line_items:
line_items_json.append({
"name": item.get("name", ""),
"description": item.get("description", ""),
"value": item.get("value"),
"quantity": item.get("quantity"),
"unit": item.get("unit", ""),
})
# Parse tax
tax_data = quote.get("tax", {}) or {}
tax_amount = tax_data.get("total")
# Parse discount
discount_data = quote.get("discount", {}) or {}
discount_amount = discount_data.get("total")
existing = JobberQuote.query.filter_by(
company_id=self.company_id,
external_id=external_id,
).first()
# Map Jobber quote status to our funnel stage
jobber_status = quote.get("status", "").lower()
status_map = {
"pending": "pending",
"draft": "draft",
"accepted": "accepted",
"rejected": "rejected",
"converted": "converted",
}
mapped_status = status_map.get(jobber_status, jobber_status)
quote_data = {
"external_id": external_id,
"name": quote.get("name", "") or "",
"status": mapped_status,
"amount": quote.get("quoteValue"),
"currency": quote.get("currencyCode", "USD") or "USD",
"tax_amount": tax_amount,
"discount_amount": discount_amount,
"sent_at": datetime.fromisoformat(quote["sentAt"].replace("Z", "+00:00")) if quote.get("sentAt") else None,
"accepted_at": datetime.fromisoformat(quote["acceptedAt"].replace("Z", "+00:00")) if quote.get("acceptedAt") else None,
"rejected_at": datetime.fromisoformat(quote["rejectedAt"].replace("Z", "+00:00")) if quote.get("rejectedAt") else None,
"expires_at": datetime.fromisoformat(quote["expiresAt"].replace("Z", "+00:00")) if quote.get("expiresAt") else None,
"line_items_json": line_items_json,
"properties_json": {"raw": quote},
"synced_at": now,
}
if existing:
existing.__dict__.update(quote_data)
else:
quote_data["company_id"] = self.company_id
new = JobberQuote(**quote_data)
db.session.add(new)
db.session.flush()
# Also create an Estimate record for funnel tracking
self._create_estimate_from_quote(new, client_external)
db.session.commit()
def _create_estimate_from_quote(self, quote: "JobberQuote", client_external_id: str):
"""Create an Estimate record from a Jobber quote for funnel tracking."""
from app.models import db, Estimate, JobberClient
client = JobberClient.query.filter_by(
company_id=self.company_id,
external_id=client_external_id,
).first()
# Map status to funnel stage
stage_map = {
"draft": "draft",
"pending": "delivered",
"accepted": "accepted",
"rejected": "rejected",
"converted": "closed",
}
now = datetime.now(timezone.utc)
funnel_stage = stage_map.get(quote.status, "draft")
# Build stage history
stage_history = [{"stage": funnel_stage, "timestamp": now.isoformat()}]
estimate = Estimate(
company_id=self.company_id,
estimate_number=quote.name or f"JB-{quote.external_id[:8]}",
total_value=quote.amount,
stage=funnel_stage,
stage_history=stage_history,
stage_changed_at=now,
prospect_name=client.full_name if client else quote.name,
prospect_email=client.email if client else "",
prospect_phone=client.phone if client else "",
description=quote.description,
created_at=datetime.fromisoformat(quote.created_at.isoformat()) if quote.created_at else now,
updated_at=now,
)
if quote.sent_at:
estimate.delivered_at = quote.sent_at
if quote.accepted_at:
estimate.accepted_at = quote.accepted_at
if quote.expired_at:
estimate.expires_at = quote.expired_at
db.session.add(estimate)
def _sync_jobs(self):
"""Sync jobs from Jobber."""
raw_jobs = self._paginated_query(JOBS_QUERY, "jobs")
logger.info("Jobber: fetched %d jobs", len(raw_jobs))
from app.models import db, JobberJob
now = datetime.now(timezone.utc)
for job in raw_jobs:
external_id = job.get("id", "")
if not external_id:
continue
address = job.get("address", {}) or {}
value_data = job.get("jobValue", {}) or {}
cost_data = job.get("jobCost", {}) or {}
def parse_dt(val):
if val:
return datetime.fromisoformat(val.replace("Z", "+00:00"))
return None
job_data = {
"external_id": external_id,
"name": job.get("name", "") or "",
"description": job.get("description", "") or "",
"status": job.get("status", "draft").lower() or "draft",
"project_type": job.get("type", "") or "",
"estimated_revenue": value_data.get("estimated"),
"actual_revenue": value_data.get("actual"),
"estimated_cost": cost_data.get("estimated"),
"actual_cost": cost_data.get("actual"),
"start_date": parse_dt(job.get("scheduledStart")),
"estimated_end_date": parse_dt(job.get("scheduledEnd")),
"end_date": parse_dt(job.get("actualEnd")),
"address_line1": address.get("line1", "") or "",
"address_line2": address.get("line2", "") or "",
"city": address.get("city", "") or "",
"state": address.get("province", "") or "",
"postal_code": address.get("postalCode", "") or "",
"properties_json": {"raw": job},
"synced_at": now,
}
existing = JobberJob.query.filter_by(
company_id=self.company_id,
external_id=external_id,
).first()
if existing:
existing.__dict__.update(job_data)
else:
job_data["company_id"] = self.company_id
new = JobberJob(**job_data)
db.session.add(new)
db.session.commit()
def _sync_invoices(self):
"""Sync invoices from Jobber."""
raw_invoices = self._paginated_query(INVOICES_QUERY, "invoices")
logger.info("Jobber: fetched %d invoices", len(raw_invoices))
from app.models import db, JobberInvoice
now = datetime.now(timezone.utc)
def parse_dt(val):
if val:
return datetime.fromisoformat(val.replace("Z", "+00:00"))
return None
for inv in raw_invoices:
external_id = inv.get("id", "")
if not external_id:
continue
line_items = inv.get("lineItems", []) or []
line_items_json = [
{
"name": item.get("name", ""),
"quantity": item.get("quantity"),
"unit": item.get("unit", ""),
"value": item.get("value"),
}
for item in line_items
]
job_data = inv.get("job", {}) or {}
inv_data = {
"external_id": external_id,
"invoice_number": inv.get("number", "") or "",
"status": inv.get("status", "draft").lower() or "draft",
"amount": inv.get("total"),
"currency": inv.get("currencyCode", "USD") or "USD",
"tax_amount": inv.get("taxTotal"),
"paid_amount": inv.get("paidTotal", 0.0) or 0.0,
"issued_at": parse_dt(inv.get("createdAt")),
"due_date": parse_dt(inv.get("dueDate")),
"paid_at": parse_dt(inv.get("paidAt")),
"line_items_json": line_items_json,
"properties_json": {"raw": inv},
"synced_at": now,
}
existing = JobberInvoice.query.filter_by(
company_id=self.company_id,
external_id=external_id,
).first()
if existing:
existing.__dict__.update(inv_data)
else:
inv_data["company_id"] = self.company_id
new = JobberInvoice(**inv_data)
db.session.add(new)
db.session.commit()
# -- status --------------------------------------------------------------
def status(self) -> Dict[str, Any]:
"""Check connection health and record counts."""
try:
from app.models import db, JobberClient, JobberQuote, JobberJob, JobberInvoice
client_count = JobberClient.query.filter_by(
company_id=self.company_id
).count()
has_token = bool(self._access_token or self.config.get("access_token"))
return {
"status": "connected" if self._connected and has_token else "disconnected",
"connected": self._connected,
"has_token": has_token,
"record_counts": {
"clients": client_count,
"quotes": JobberQuote.query.filter_by(company_id=self.company_id).count(),
"jobs": JobberJob.query.filter_by(company_id=self.company_id).count(),
"invoices": JobberInvoice.query.filter_by(company_id=self.company_id).count(),
},
"token_expires_at": self._token_expires_at.isoformat() if self._token_expires_at else None,
}
except Exception as exc:
return {
"status": "error",
"error": str(exc),
}
# -- OAuth methods -------------------------------------------------------
def oauth_authorize_url(self, state: str) -> str:
"""Build Jobber authorization URL."""
client_id = self._get_client_id()
redirect_uri = self._get_redirect_uri()
scopes = "+".join(self.OAUTH_SCOPES)
return (
f"{_AUTHORIZE_URL}?"
f"response_type=code&"
f"client_id={client_id}&"
f"redirect_uri={requests.utils.quote(redirect_uri, safe='')}&"
f"scope={scopes}&"
f"state={state}"
)
def exchange_code_for_tokens(self, code: str, **kwargs) -> Dict[str, Any]:
"""Exchange authorization code for tokens."""
client_id = self._get_client_id()
client_secret = self._get_client_secret()
redirect_uri = self._get_redirect_uri()
resp = requests.post(
_TOKEN_URL,
data={
"grant_type": "authorization_code",
"code": code,
"client_id": client_id,
"client_secret": client_secret,
"redirect_uri": redirect_uri,
},
timeout=30,
)
resp.raise_for_status()
data = resp.json()
self._access_token = data.get("access_token")
self._refresh_token = data.get("refresh_token")
expires_in = int(data.get("expires_in", 3600))
self._token_expires_at = datetime.now(timezone.utc) + timedelta(seconds=expires_in)
# Persist in config
self.config["access_token"] = self._access_token
self.config["refresh_token"] = self._refresh_token
logger.info("Jobber OAuth tokens exchanged successfully")
return {
"access_token": self._access_token,
"refresh_token": self._refresh_token,
"expires_in": expires_in,
}
# -- Register in the connector registry --------------------------------------
_REGISTRY["jobber"] = JobberConnector