NC runtime: scope activity checks per connection user

This commit is contained in:
Terranom674
2026-08-18 19:07:20 +02:00
parent b301853ea2
commit 69815b1cf4

View File

@@ -1,5 +1,5 @@
#!/usr/bin/env python3 #!/usr/bin/env python3
"""Debounce Nextcloud activity and request a periodic safety reconciliation.""" """Debounce Nextcloud activity for one connector scope."""
from __future__ import annotations from __future__ import annotations
@@ -24,54 +24,117 @@ def validate_view(name: str) -> str:
return value return value
def sql_literal(value: str) -> str:
return "'" + str(value).replace("'", "''") + "'"
def load(path: Path) -> dict[str, object]: def load(path: Path) -> dict[str, object]:
defaults: dict[str, object] = {"processed":0,"observed":0,"pending_since":0,"last_change":0,"last_full":0,"source_signature":""} defaults: dict[str, object] = {
"processed": 0, "observed": 0, "pending_since": 0,
"last_change": 0, "last_full": 0, "source_signature": "",
}
if not path.exists(): if not path.exists():
return defaults return defaults
data = json.loads(path.read_text(encoding="utf-8")) data = json.loads(path.read_text(encoding="utf-8"))
return {"processed":int(data.get("processed",0)),"observed":int(data.get("observed",0)),"pending_since":int(data.get("pending_since",0)),"last_change":int(data.get("last_change",0)),"last_full":int(data.get("last_full",0)),"source_signature":str(data.get("source_signature",""))} return {
"processed": int(data.get("processed", 0)),
"observed": int(data.get("observed", 0)),
"pending_since": int(data.get("pending_since", 0)),
"last_change": int(data.get("last_change", 0)),
"last_full": int(data.get("last_full", 0)),
"source_signature": str(data.get("source_signature", "")),
}
def save(path: Path, state: dict[str, object]) -> None: def save(path: Path, state: dict[str, object]) -> None:
path.parent.mkdir(parents=True, exist_ok=True) path.parent.mkdir(parents=True, exist_ok=True)
fd, name = tempfile.mkstemp(dir=path.parent) fd, name = tempfile.mkstemp(dir=path.parent)
with os.fdopen(fd, "w", encoding="utf-8") as handle: with os.fdopen(fd, "w", encoding="utf-8") as handle:
json.dump(state,handle,sort_keys=True);handle.write("\n") json.dump(state, handle, sort_keys=True)
handle.write("\n")
Path(name).replace(path) Path(name).replace(path)
def query(args: argparse.Namespace, sql: str) -> str: def query(args: argparse.Namespace, sql: str) -> str:
env=os.environ.copy();env["PGPASSWORD"]=args.password_file.read_text(encoding="utf-8").strip() env = os.environ.copy()
command=["psql","-XAt","-h",args.host,"-p",str(args.port),"-U",args.user,"-d",args.database,"-v","ON_ERROR_STOP=1","-c",sql] env["PGPASSWORD"] = args.password_file.read_text(encoding="utf-8").strip()
command = [
"psql", "-XAt", "-h", args.host, "-p", str(args.port),
"-U", args.user, "-d", args.database, "-v", "ON_ERROR_STOP=1", "-c", sql,
]
return subprocess.run(command, env=env, check=True, text=True, capture_output=True).stdout return subprocess.run(command, env=env, check=True, text=True, capture_output=True).stdout
def user_where(args: argparse.Namespace) -> str:
if not args.access_user:
return ""
return " WHERE lower(access_user) = lower(" + sql_literal(args.access_user) + ")"
def latest(args: argparse.Namespace) -> int: def latest(args: argparse.Namespace) -> int:
return int(query(args,f"SELECT COALESCE(MAX(activity_id), 0) FROM {validate_view(args.view)}").strip()) view = validate_view(args.view)
return int(query(args, f"SELECT COALESCE(MAX(activity_id), 0) FROM {view}{user_where(args)}").strip())
def source_signature(args: argparse.Namespace) -> str: def source_signature(args: argparse.Namespace) -> str:
if not args.source_view:return "" if args.roots_config:
payload=query(args,f"SELECT share_id, display_name, storage_id, source_path FROM {validate_view(args.source_view)} ORDER BY share_id") return hashlib.sha256(args.roots_config.read_bytes()).hexdigest()
if not args.source_view:
return ""
view = validate_view(args.source_view)
payload = query(args, f"SELECT share_id, display_name, storage_id, source_path FROM {view}{user_where(args)} ORDER BY share_id")
return hashlib.sha256(payload.encode("utf-8")).hexdigest() return hashlib.sha256(payload.encode("utf-8")).hexdigest()
def main() -> int: def main() -> int:
parser=argparse.ArgumentParser();parser.add_argument("action",choices=("check","commit"));parser.add_argument("--state",required=True,type=Path);parser.add_argument("--host",required=True);parser.add_argument("--port",type=int,default=5432);parser.add_argument("--database",required=True);parser.add_argument("--user",required=True);parser.add_argument("--password-file",required=True,type=Path);parser.add_argument("--view",default="piwigo_showcase_activity");parser.add_argument("--source-view",default="");parser.add_argument("--quiet",type=int,default=120);parser.add_argument("--max-wait",type=int,default=900);parser.add_argument("--full-after",type=int,default=86400);args=parser.parse_args() parser = argparse.ArgumentParser()
parser.add_argument("action", choices=("check", "commit"))
parser.add_argument("--state", required=True, type=Path)
parser.add_argument("--host", required=True)
parser.add_argument("--port", type=int, default=5432)
parser.add_argument("--database", required=True)
parser.add_argument("--user", required=True)
parser.add_argument("--password-file", required=True, type=Path)
parser.add_argument("--view", required=True)
parser.add_argument("--source-view", default="")
parser.add_argument("--roots-config", type=Path)
parser.add_argument("--access-user", default="")
parser.add_argument("--quiet", type=int, default=120)
parser.add_argument("--max-wait", type=int, default=900)
parser.add_argument("--full-after", type=int, default=86400)
args = parser.parse_args()
try: try:
validate_view(args.view) validate_view(args.view)
if args.source_view:validate_view(args.source_view) if args.source_view:
state=load(args.state);now=int(time.time());current=latest(args);current_source=source_signature(args) validate_view(args.source_view)
if args.action=="commit":state.update(processed=current,observed=current,pending_since=0,last_change=0,last_full=now,source_signature=current_source);save(args.state,state);return 0 state = load(args.state)
if not int(state["last_full"]):return 0 now = int(time.time())
if current_source and current_source!=str(state["source_signature"]):return 0 current = latest(args)
if current>int(state["observed"]):state["observed"]=current;state["last_change"]=now;state["pending_since"]=int(state["pending_since"]) or now;save(args.state,state) current_source = source_signature(args)
if now-int(state["last_full"])>=args.full_after:return 0 if args.action == "commit":
if int(state["observed"])<=int(state["processed"]):return 3 state.update(processed=current, observed=current, pending_since=0, last_change=0, last_full=now, source_signature=current_source)
if now-int(state["last_change"])>=args.quiet or now-int(state["pending_since"])>=args.max_wait:return 0 save(args.state, state)
return 0
if not int(state["last_full"]):
return 0
if current_source and current_source != str(state["source_signature"]):
return 0
if current > int(state["observed"]):
state["observed"] = current
state["last_change"] = now
state["pending_since"] = int(state["pending_since"]) or now
save(args.state, state)
if now - int(state["last_full"]) >= args.full_after:
return 0
if int(state["observed"]) <= int(state["processed"]):
return 3
if now - int(state["last_change"]) >= args.quiet or now - int(state["pending_since"]) >= args.max_wait:
return 0
return 3 return 3
except Exception as error: except Exception as error:
print(f"activity-gate: {error}",file=sys.stderr);return 1 print(f"activity-gate: {error}", file=sys.stderr)
return 1
if __name__=="__main__":raise SystemExit(main()) if __name__ == "__main__":
raise SystemExit(main())