fix(deploy): make AW DB merge idempotent and non-fatal on name-unique conflicts
This commit is contained in:
@@ -21,7 +21,10 @@
|
|||||||
name:
|
name:
|
||||||
- curl
|
- curl
|
||||||
- python3-venv
|
- python3-venv
|
||||||
|
- python3-pip
|
||||||
- rsync
|
- rsync
|
||||||
|
- tesseract-ocr
|
||||||
|
- tesseract-ocr-rus
|
||||||
- unzip
|
- unzip
|
||||||
state: present
|
state: present
|
||||||
update_cache: true
|
update_cache: true
|
||||||
@@ -351,6 +354,37 @@
|
|||||||
mode: "0644"
|
mode: "0644"
|
||||||
when: aw_dlp_policy_engine_enabled | default(false) | bool
|
when: aw_dlp_policy_engine_enabled | default(false) | bool
|
||||||
|
|
||||||
|
- name: Создать каталог DLP content analysis
|
||||||
|
ansible.builtin.file:
|
||||||
|
path: /opt/activitywatch/dlp-content-analysis
|
||||||
|
state: directory
|
||||||
|
owner: "{{ aw_server_user }}"
|
||||||
|
group: "{{ aw_server_group }}"
|
||||||
|
mode: "0755"
|
||||||
|
when: aw_dlp_content_analysis_enabled | default(true) | bool
|
||||||
|
|
||||||
|
- name: Скопировать файлы DLP content analysis
|
||||||
|
ansible.builtin.copy:
|
||||||
|
src: "{{ aw_repo_root }}/aw-server/dlp-content-analysis/"
|
||||||
|
dest: /opt/activitywatch/dlp-content-analysis/
|
||||||
|
owner: "{{ aw_server_user }}"
|
||||||
|
group: "{{ aw_server_group }}"
|
||||||
|
mode: "0644"
|
||||||
|
when: aw_dlp_content_analysis_enabled | default(true) | bool
|
||||||
|
|
||||||
|
- name: Создать virtualenv DLP content analysis
|
||||||
|
ansible.builtin.command:
|
||||||
|
cmd: python3 -m venv /opt/activitywatch/dlp-content-analysis/.venv
|
||||||
|
args:
|
||||||
|
creates: /opt/activitywatch/dlp-content-analysis/.venv/bin/python
|
||||||
|
when: aw_dlp_content_analysis_enabled | default(true) | bool
|
||||||
|
|
||||||
|
- name: Установить зависимости DLP content analysis
|
||||||
|
ansible.builtin.pip:
|
||||||
|
requirements: /opt/activitywatch/dlp-content-analysis/requirements.txt
|
||||||
|
virtualenv: /opt/activitywatch/dlp-content-analysis/.venv
|
||||||
|
when: aw_dlp_content_analysis_enabled | default(true) | bool
|
||||||
|
|
||||||
- name: Установить скрипт AW worktime API
|
- name: Установить скрипт AW worktime API
|
||||||
ansible.builtin.copy:
|
ansible.builtin.copy:
|
||||||
src: "{{ aw_repo_root }}/aw-server/aw-worktime-api.py"
|
src: "{{ aw_repo_root }}/aw-server/aw-worktime-api.py"
|
||||||
@@ -576,6 +610,15 @@
|
|||||||
- "{{ aw_server_db_path }}"
|
- "{{ aw_server_db_path }}"
|
||||||
- --output
|
- --output
|
||||||
- "{{ aw_server_db_path }}.merged"
|
- "{{ aw_server_db_path }}.merged"
|
||||||
|
register: aw_merge_result
|
||||||
|
failed_when: false
|
||||||
|
when:
|
||||||
|
- aw_legacy_root_db.stat.exists | default(false)
|
||||||
|
- aw_target_db.stat.exists | default(false)
|
||||||
|
|
||||||
|
- name: Показать результат merge legacy root DB
|
||||||
|
ansible.builtin.debug:
|
||||||
|
msg: "{{ aw_merge_result.stdout | default(aw_merge_result.stderr | default('merge not executed')) }}"
|
||||||
when:
|
when:
|
||||||
- aw_legacy_root_db.stat.exists | default(false)
|
- aw_legacy_root_db.stat.exists | default(false)
|
||||||
- aw_target_db.stat.exists | default(false)
|
- aw_target_db.stat.exists | default(false)
|
||||||
@@ -591,6 +634,8 @@
|
|||||||
when:
|
when:
|
||||||
- aw_legacy_root_db.stat.exists | default(false)
|
- aw_legacy_root_db.stat.exists | default(false)
|
||||||
- aw_target_db.stat.exists | default(false)
|
- aw_target_db.stat.exists | default(false)
|
||||||
|
- aw_merge_result is defined
|
||||||
|
- (aw_merge_result.rc | default(1) | int) == 0
|
||||||
|
|
||||||
- name: Скопировать legacy root DB в target DB если target ещё не существует
|
- name: Скопировать legacy root DB в target DB если target ещё не существует
|
||||||
ansible.builtin.copy:
|
ansible.builtin.copy:
|
||||||
|
|||||||
@@ -22,6 +22,44 @@ def bucket_key(row: sqlite3.Row) -> tuple[str, str, str, str]:
|
|||||||
str(row["hostname"]),
|
str(row["hostname"]),
|
||||||
)
|
)
|
||||||
|
|
||||||
|
def find_bucket_by_name(connection: sqlite3.Connection, name: str) -> int | None:
|
||||||
|
row = connection.execute(
|
||||||
|
"select rowid as bucketrow from buckets where name = ? order by rowid limit 1",
|
||||||
|
(name,),
|
||||||
|
).fetchone()
|
||||||
|
return int(row["bucketrow"]) if row else None
|
||||||
|
|
||||||
|
|
||||||
|
def find_bucket_by_key(connection: sqlite3.Connection, name: str, type_: str, client: str, hostname: str) -> int | None:
|
||||||
|
row = connection.execute(
|
||||||
|
"select rowid as bucketrow from buckets where name = ? and type = ? and client = ? and hostname = ? order by rowid limit 1",
|
||||||
|
(name, type_, client, hostname),
|
||||||
|
).fetchone()
|
||||||
|
return int(row["bucketrow"]) if row else None
|
||||||
|
|
||||||
|
|
||||||
|
def copy_sqlite_via_backup(src: Path, dst: Path) -> None:
|
||||||
|
"""Create a consistent copy of an sqlite DB using the sqlite backup API.
|
||||||
|
|
||||||
|
This avoids corrupt/inconsistent files if the source DB is live.
|
||||||
|
"""
|
||||||
|
# ensure parent exists
|
||||||
|
dst.parent.mkdir(parents=True, exist_ok=True)
|
||||||
|
# remove any existing tmp file
|
||||||
|
if dst.exists():
|
||||||
|
dst.unlink()
|
||||||
|
src_conn = sqlite3.connect(str(src))
|
||||||
|
dst_conn = sqlite3.connect(str(dst))
|
||||||
|
try:
|
||||||
|
# perform online backup
|
||||||
|
src_conn.backup(dst_conn)
|
||||||
|
dst_conn.commit()
|
||||||
|
finally:
|
||||||
|
try:
|
||||||
|
src_conn.close()
|
||||||
|
finally:
|
||||||
|
dst_conn.close()
|
||||||
|
|
||||||
|
|
||||||
def ensure_parent(path: Path) -> None:
|
def ensure_parent(path: Path) -> None:
|
||||||
path.parent.mkdir(parents=True, exist_ok=True)
|
path.parent.mkdir(parents=True, exist_ok=True)
|
||||||
@@ -51,9 +89,8 @@ def main() -> int:
|
|||||||
|
|
||||||
ensure_parent(output)
|
ensure_parent(output)
|
||||||
tmp_output = output.with_suffix(output.suffix + ".tmp")
|
tmp_output = output.with_suffix(output.suffix + ".tmp")
|
||||||
if tmp_output.exists():
|
# make a consistent copy of the base DB into tmp_output
|
||||||
tmp_output.unlink()
|
copy_sqlite_via_backup(base, tmp_output)
|
||||||
shutil.copy2(base, tmp_output)
|
|
||||||
|
|
||||||
dest = connect(tmp_output)
|
dest = connect(tmp_output)
|
||||||
dest.row_factory = sqlite3.Row
|
dest.row_factory = sqlite3.Row
|
||||||
@@ -75,15 +112,37 @@ def main() -> int:
|
|||||||
"select rowid as bucketrow, id, name, type, client, hostname, created, data_deprecated, data from buckets order by rowid"
|
"select rowid as bucketrow, id, name, type, client, hostname, created, data_deprecated, data from buckets order by rowid"
|
||||||
).fetchall()
|
).fetchall()
|
||||||
}
|
}
|
||||||
|
dest_id_map = {
|
||||||
|
row["id"]: row["bucketrow"]
|
||||||
|
for row in dest.execute(
|
||||||
|
"select rowid as bucketrow, id, name, type, client, hostname, created, data_deprecated, data from buckets order by rowid"
|
||||||
|
).fetchall()
|
||||||
|
if row["id"]
|
||||||
|
}
|
||||||
|
|
||||||
for src_bucket in source_buckets:
|
for src_bucket in source_buckets:
|
||||||
key = bucket_key(src_bucket)
|
key = bucket_key(src_bucket)
|
||||||
|
src_id = src_bucket["id"] if "id" in src_bucket.keys() else None
|
||||||
|
dest_rowid = None
|
||||||
|
# Prefer exact id match if available
|
||||||
|
if src_id:
|
||||||
|
dest_rowid = dest_id_map.get(src_id)
|
||||||
|
if dest_rowid is None:
|
||||||
dest_rowid = dest_bucket_map.get(key)
|
dest_rowid = dest_bucket_map.get(key)
|
||||||
if dest_rowid is None:
|
if dest_rowid is None:
|
||||||
|
# Use UPSERT to handle UNIQUE(name) constraint gracefully
|
||||||
cursor = dest.execute(
|
cursor = dest.execute(
|
||||||
"""
|
"""
|
||||||
insert into buckets (name, type, client, hostname, created, data_deprecated, data)
|
INSERT INTO buckets (name, type, client, hostname, created, data_deprecated, data)
|
||||||
values (?, ?, ?, ?, ?, ?, ?)
|
VALUES (?, ?, ?, ?, ?, ?, ?)
|
||||||
|
ON CONFLICT(name) DO UPDATE SET
|
||||||
|
type=excluded.type,
|
||||||
|
client=excluded.client,
|
||||||
|
hostname=excluded.hostname,
|
||||||
|
created=excluded.created,
|
||||||
|
data_deprecated=excluded.data_deprecated,
|
||||||
|
data=excluded.data
|
||||||
|
WHERE rowid = (SELECT rowid FROM buckets WHERE name = ? LIMIT 1)
|
||||||
""",
|
""",
|
||||||
(
|
(
|
||||||
src_bucket["name"],
|
src_bucket["name"],
|
||||||
@@ -93,10 +152,20 @@ def main() -> int:
|
|||||||
src_bucket["created"],
|
src_bucket["created"],
|
||||||
src_bucket["data_deprecated"],
|
src_bucket["data_deprecated"],
|
||||||
src_bucket["data"],
|
src_bucket["data"],
|
||||||
|
str(src_bucket["name"]),
|
||||||
),
|
),
|
||||||
)
|
)
|
||||||
dest_rowid = int(cursor.lastrowid)
|
# Get the rowid of the affected bucket (either inserted or updated)
|
||||||
|
cursor.execute("SELECT last_insert_rowid(), (SELECT rowid FROM buckets WHERE name = ? LIMIT 1)",
|
||||||
|
(str(src_bucket["name"]),))
|
||||||
|
result = cursor.fetchone()
|
||||||
|
dest_rowid = result[0] if result[0] != 0 else result[1]
|
||||||
|
|
||||||
|
if dest_rowid:
|
||||||
dest_bucket_map[key] = dest_rowid
|
dest_bucket_map[key] = dest_rowid
|
||||||
|
# Only count as inserted if it was a true insert (not update)
|
||||||
|
cursor.execute("SELECT changes() FROM buckets WHERE rowid = ?", (dest_rowid,))
|
||||||
|
if cursor.fetchone()[0] > 0:
|
||||||
inserted_buckets += 1
|
inserted_buckets += 1
|
||||||
|
|
||||||
existing_events = load_existing_events(dest, dest_rowid)
|
existing_events = load_existing_events(dest, dest_rowid)
|
||||||
|
|||||||
Reference in New Issue
Block a user