Files
content-ingestion-agent/kab_ingestion/sharepoint.py

27 lines
2.5 KiB
Python

from .models import *
from .ports import HTTP, SecretStore
from .util import *
class SharePointConnector:
"""Microsoft Graph drive connector for PDF, DOCX, and HTML files."""
def __init__(self,http:HTTP,secrets:SecretStore,clock=now): self.http,self.secrets,self.clock=http,secrets,clock
def _url(self,c,path): return "https://graph.microsoft.com/v1.0/sites/"+c.scope["site_id"]+path
def _doc(self,c,item,token,event=None):
content=self.http.request("GET",item["@microsoft.graph.downloadUrl"],headers={}).get("content","")
ext=item["name"].lower().rsplit('.',1)[-1]; mime={"pdf":"application/pdf","docx":"application/vnd.openxmlformats-officedocument.wordprocessingml.document","html":"text/html","htm":"text/html"}[ext]
p=Provenance("sharepoint",item["id"],item.get("webUrl",""),item.get("eTag"),self.clock(),event)
acl=ACL(c.tenant_id,tuple(c.options.get("principals",[])),tuple(c.options.get("groups",[])))
return Document(stable_id(c.tenant_id,"sharepoint",item["id"]),c.tenant_id,item["name"],content,mime,digest(content),item.get("lastModifiedDateTime"),p,acl,{"library_id":c.scope.get("library_id"),"folder":c.scope.get("folder")})
def full(self,c):
require_scope(c.scope,("site_id","library_id")); token=self.secrets.get(c.secret_ref); path=f"/drives/{c.scope['library_id']}/root/children"
data=self.http.request("GET",self._url(c,path),headers=json_headers(token)); wanted=tuple(c.scope.get("extensions",[".pdf",".docx",".html"]))
docs=[self._doc(c,x,token) for x in data.get("value",[]) if x.get("file") and x["name"].lower().endswith(wanted) and (not c.scope.get("folder") or x.get("parentReference",{}).get("path","").endswith(c.scope["folder"]))]
return docs,SyncCursor(c.connector_id,"full",delta_token=data.get("@odata.deltaLink"),updated_at=self.clock())
def incremental(self,c,cursor,changes=None):
token=self.secrets.get(c.secret_ref); data=self.http.request("GET",cursor.delta_token or self._url(c,f"/drives/{c.scope['library_id']}/root/delta"),headers=json_headers(token)); docs=[]; deleted=[]
for x in data.get("value",[]):
if "deleted" in x: deleted.append(x["id"])
elif x.get("file"): docs.append(self._doc(c,x,token))
return docs,SyncCursor(c.connector_id,"incremental",delta_token=data.get("@odata.deltaLink",cursor.delta_token),updated_at=self.clock())
def webhook(self,c,payload,headers): return [Change(x.get("resource",""),"modified",x.get("sequenceNumber"),payload.get("subscriptionId")) for x in payload.get("value",[])]