Notice of Pre-AIA or AIA Status
The present application, filed on or after March 16, 2013, is being examined under the first inventor to file provisions of the AIA .
Response to Amendment
In response to the amendment filed on June 3, 2026:
Claims 1-20 are canceled.
Claims 28, 36, and 38-40 are amended.
Claims 21-40 are pending.
Response to Arguments
In response to the remarks filed on June 3, 2026:
a. The double patenting rejections of the pending claims are maintained since Applicant wished to resolve the rejections at a later time.
b. Applicant’s remarks regarding 35 U.S.C. 101 rejections of claims 21-22, and 24-40 the pending claims have been fully considered but are not persuasive.
Step 2A – Prong 1 (Judicial Exception Grouping)
Applicant remarks that each of independent claims 21, 36, and 39 requires operations across physical storage infrastructure and cannot practically be performed in the human mind or with pen/paper (SRI Int’l v Cisco).
The Examiner respectfully disagrees. The USPTO 2019/2024 Revised Patent Subject Matter Eligibility Guidance clarifies that an abstract idea does not require a claim to be exclusively a mental process. It can also fall under Certain Methods of Organizing Human Activity (managing data records, data synchronization, and data architecture). The recitation of physical storage infrastructure of plurality of source/target storage systems is merely the environment in which the abstract process of data record synchronization is performed. Under Electric Power Group, LLC v. Alstom S.A., 830 F.3d 1350 (Fed. Cir. 2016), collecting, analyzing, and transmitting data points across a network remains an abstract concept even if executed on networked hardware.
Step 2A – Prong 2 (Integration into a Practical Application)
Applicant remarks that each of the independent claims addresses a specific technical problem (causal consistency in distributed storage) and improve computer functionality (Amdocs v . Openet Telecom).
The Examiner respectfully disagrees. Applicant relies on technical capabilities described in the specification (e.g., enforcing dependency relationships, causal consistency) that are not actually recited in the language of the independent claims. As drafted, each independent claim merely recites high-level functional results, i.e., replicating…local source checkpoints; and creating…a coordinated target checkpoint. Each claim fails to recite how the system programmatic structures, algorithms, or storage controllers achieve this state consistency. Under MPEP 2106.04(d), an alleged technical improvement must be reflected in the recited claim limitations, not merely stated in the remarks or spec. Generic data copying and index creation across conventional storage hardware do not integrate the abstract concept into a practical application.
Step 2B (Inventive Concept/Providing Significantly More)
Applicant remarks that the ordered combination of steps is unconventional and the Examiner failed to provided evidence of well-understood, routine, and conventional (WURC) activity.
The Examiner respectfully disagrees. The claim components are conventional network data storage units executing routine data replication and state file creation operations. Combining routine data replication with routine checkpoint generation is a logical aggregation of conventional steps. Without specific technical limitations detailing how distributed consensus or non-block snapshots are executed, the ordered combination yields nothing more than the sum of its generic parts. As such, claims 21-22, and 24-40 are ineligible under 35 U.S.C. 101.
c. Applicant’s remarks regarding 35 U.S.C. 103 rejections of claims 21, 24-28, 30-40 the pending claims have been fully considered but are not persuasive.
(i) Regarding independent claims 21, 36, and 39, Applicant remarks Flynn discloses a single checkpoint generated centrally by engine 414 for a single application instance, and storage devices A-C are merely storage locations within one production site while Annamalai lacks checkpoint generation and coordination, and there is no motivation to combine Annamalai with Flynn.
The Examiner respectfully disagrees. Each independent claim recites a plurality of source storage systems and Flynn explicitly describes a production site 450 comprising distinct physical or logical storage devices A, B, and C. The disclosure of Flynn shows that these devices can be housed within a single site or separated across networks does not restrict them from being systems under Broadest Reasonable Interpretation (BRI) of the claimed source storage systems. The fact that paragraph [0111] of Flynn mentions an alternative embodiment where the site comprises a single device does not negate or erase the primary disclosure of multiple storage units since Flynn explicitly indicates “Rather, any number of storage devices may be utilized without departing from the spirit and scope of the present invention” in the same paragraph.
In Flynn, application data copy is stored in storage device A while stateful metadata is stored in storage device B ([0093]-[0095], [0097], and [0099]-[0104]). These state copies generated at substantially the same time ([0070]) functionally form a local point-in-time snapshots (local checkpoints) of their respective storage components. Applicant appears to import an unrecited “independent generation” requirement into the independent claims. Each independent claim merely requires that local checkpoints are associated with a coordinated checkpoint, which Flynn’s engine achieves across at least devices A and B.
Further, Applicant is noted that under 35 U.S.C. 101, references must be evaluated as a combined whole (In re Merck & Co.) Flynn teaches coordinating state checkpoints across storage volumes, while Annamalai teach partitioning datasets into distributed shards across a plurality of network servers. It would have been obvious to a person skilled in the art (PHOSITA) at the time of the invention to apply Flynn’s check point coordination logic to Annamalai’s distributed sharded dataset to ensure point-in-time consistency across large scale distributed storage nodes. The combination yields predictable results under KSR Int’l Co. v. Teleflex Inc.
(ii) Applicant’s remarks regarding the 35 U.S.C. 103 rejection of claim 28 have been fully considered but are moot in view of a new ground of rejection presented here on.
(iii) 35 U.S.C. 103 rejections of dependent claims 24-27, 30, 32-35, 37-38, and 40 are maintained for the same reasons presented above.
c. 35 U.S.C. 103 rejections of claims 22, and 29 are withdrawn.
d. 35 U.S.C. 101 and 103 rejections of claim 23 are withdrawn. The claim is considered allowable if fully incorporate into independent claim 21.
e. Applicant’s remarks regarding the 35 U.S.C. 103 rejection of claim 31 have been fully considered but are not persuasive.
Applicant indicates that claim 31 recites identifying an alternate target system that has already received a missing local checkpoint and incorporating it into the target checkpoint and argues that Annamalai’s server replacement is generic infrastructure failover (deploying a new server), not searching for a node that already holds the specific missing checkpoint state.
The Examiner respectfully disagrees. Annamalai discloses configurable fault tolerance within sync replica sets where data is replicated across a hierarchy ([0028]). When a primary or secondary node fails, the system evaluates the state of remaining nodes in the sync replica sets to identify which follower already possesses the persisted log/version key up to the required sequence number. Re-routing the target snapshot synthesis to a surviving replica node that already holds the synced log is the natural operation of Annamalai’s quorum synchronization mechanism.
Double Patenting
The nonstatutory double patenting rejection is based on a judicially created doctrine grounded in public policy (a policy reflected in the statute) so as to prevent the unjustified or improper timewise extension of the “right to exclude” granted by a patent and to prevent possible harassment by multiple assignees. A nonstatutory double patenting rejection is appropriate where the claims at issue are not identical, but at least one examined application claim is not patentably distinct from the reference claim(s) because the examined application claim is either anticipated by, or would have been obvious over, the reference claim(s). See, e.g., In re Berg, 140 F.3d 1428, 46 USPQ2d 1226 (Fed. Cir. 1998); In re Goodman, 11 F.3d 1046, 29 USPQ2d 2010 (Fed. Cir. 1993); In re Longi, 759 F.2d 887, 225 USPQ 645 (Fed. Cir. 1985); In re Van Ornum, 686 F.2d 937, 214 USPQ 761 (CCPA 1982); In re Vogel, 422 F.2d 438, 164 USPQ 619 (CCPA 1970); and In re Thorington, 418 F.2d 528, 163 USPQ 644 (CCPA 1969).
A timely filed terminal disclaimer in compliance with 37 CFR 1.321(c) or 1.321(d) may be used to overcome an actual or provisional rejection based on a nonstatutory double patenting ground provided the reference application or patent either is shown to be commonly owned with this application, or claims an invention made as a result of activities undertaken within the scope of a joint research agreement. See MPEP § 717.02 for applications subject to examination under the first inventor to file provisions of the AIA as explained in MPEP § 2159. See MPEP §§ 706.02(l)(1) - 706.02(l)(3) for applications not subject to examination under the first inventor to file provisions of the AIA . A terminal disclaimer must be signed in compliance with 37 CFR 1.321(b).
The USPTO Internet website contains terminal disclaimer forms which may be used. Please visit www.uspto.gov/forms/. The filing date of the application in which the form is filed determines what form (e.g., PTO/SB/25, PTO/SB/26, PTO/AIA /25, or PTO/AIA /26) should be used. A web-based eTerminal Disclaimer may be filled out completely online using web-screens. An eTerminal Disclaimer that meets all requirements is auto-processed and approved immediately upon submission. For more information about eTerminal Disclaimers, refer to http://www.uspto.gov/patents/process/file/efs/guidance/eTD-info-I.jsp.
Claims 21-40 are rejected on the ground of nonstatutory double patenting over claims 1-20 of Pat. No. US 12166820.
Claims 21-40 of the instant application recite similar limitations and claims 1-20 of ‘820 as being compared in the table below. For the purpose of illustration, only claims 21-35 (method claims) of the instant application are compared to the claims of the patent (underlining are used to indicate conflict limitations). The remaining claims of the instant application recite different categories (i.e., system/apparatus and medium claims) and are therefore not compared for simplicity purposes.
Instant Application
Pat. No. US 12166820
Claim 21
A method comprising:
replicating, from a plurality of source storage systems to a plurality of target storage systems, a plurality of local source checkpoints that are associated with a coordinated source checkpoint representing a snapshot of a source dataset stored across the plurality of source storage systems; and
creating, based on the replicated plurality of local source checkpoints, a coordinated target checkpoint for a replica of the source dataset that is stored across the plurality of target storage systems.
Claim 1
A method comprising:
storing a replica dataset across two or more target storage systems, wherein the replica dataset is a replication target for a source dataset stored across two or more source storage systems;
creating two or more local replicated checkpoints that are replicated from the two or more source storage systems to the two or more target storage systems, wherein two or more local source checkpoints for the two or more local replicated checkpoints are associated with a coordinated source checkpoint for the source dataset, wherein the coordinated source checkpoint represents a snapshot of the source dataset stored across the two or more source storage system; and
creating, based on the two or more local replicated checkpoints, a coordinated target checkpoint for the replica dataset.
See further Flynn and Annamalai below for mapping and motivation to combine with the claims of ‘820.
Claim 22
The method of claim 21, wherein a storage implementation of a first target storage system among the plurality of target storage systems is one of a block storage implementation, a file system implementation, a database implementation, and an object store implementation;
wherein a storage implementation of a second target storage system among the plurality of target storage systems is one of a block storage implementation, a file system implementation, a database implementation, and an object store implementation; and
wherein the storage implementation of the first target storage system is different from the storage implementation of the second target storage system.
Claim 2
The method of claim 1, wherein a storage implementation of a first target storage target system among the two or more target storage systems is one of a block storage implementation, a file system implementation, a database implementation, and an object store implementation;
wherein a storage implementation of a second target storage system among the two or more target storage systems is one of a block storage implementation, a file system implementation, a database implementation, and an object store implementation; and
wherein the storage implementation of the first target storage system is different from the storage implementation of the second target storage system.
Claim 23
The method of claim 21, wherein the plurality of target storage systems are respectively paired with the plurality of source storage systems by respective replication links to form a plurality of replicating pairs;
wherein each of the plurality of source storage systems stores a distinct local portion of the source dataset that is replicated to its paired target storage system; and
wherein each of the plurality of target storage systems stores a distinct local replicated portion of the replica.
Claim 3
The method of claim 1, wherein the two or more target storage systems are respectively paired with the two or more source storage systems by respective replication links to form a plurality of replicating pairs;
wherein each of the two or more source storage systems stores a distinct local portion of the source dataset that is replicated to its paired target storage system; and
wherein each of the target storage systems stores a distinct local replicated portion of the replica dataset.
Claim 24
The method of claim 23, wherein each local source storage system replicates its local source checkpoint to its paired target storage system.
Claim 4
The method of claim 3, wherein each local source storage system replicates its local source checkpoint to its paired target storage system
Claim 25
The method of claim 23, wherein at least one replication link employs a replication technology that is different from another replication link.
Claim 5
The method of claim 3, wherein at least one replication link employs a replication technology that is different from another replication link.
Claim 26
The method of claim 21, wherein the coordinated source checkpoint is coordinated across the plurality of source storage systems by a source coordinator.
Claim 6
The method of claim 1, wherein the coordinated source checkpoint is coordinated across the two or more source storage system by a source coordinator.
Claim 27
The method of claim 26, wherein a target coordinator receives, from the source coordinator, information relating to the coordinated source checkpoint.
Claim 7
The method of claim 6, wherein a target coordinator receives, from the source coordinator, information relating to the coordinated source checkpoint.
Claim 28
The method of claim 21, wherein plurality of source storage systems store respective distinct local portion.
Claim 6
The method of claim 1, wherein the coordinated source checkpoint is coordinated across the two or more source storage system by a source coordinator.
Claim 29
The method of claim 21, wherein a first local replicated checkpoint among two or more local replicated checkpoints of the coordinated target checkpoint is a snapshot; and
wherein a second local replicated checkpoint among the two or more local replicated checkpoints of the coordinated target checkpoint is a lightweight checkpoint.
Claim 9
The method of claim 1, wherein a first local replicated checkpoint among the two or more local replicated checkpoints of the coordinated target checkpoint is a snapshot; and
wherein a second local replicated checkpoint among the two or more local replicated checkpoints of the coordinated target checkpoint is a lightweight checkpoint.
Claim 30
The method of claim 21, wherein creating, based on two or more local replicated checkpoints, a coordinated target checkpoint for the replica includes:
determining a state of the coordinated target checkpoint.
Claim 10
The method of claim 1, wherein creating, based on the two or more local replicated checkpoints, a coordinated target checkpoint for the replica dataset includes:
determining a state of the coordinated target checkpoint.
Claim 31
The method of claim 30, wherein in response to determining a faulted state for the coordinated target checkpoint due to a particular faulted target storage system, an alternate target storage system that has received a missing local replicated checkpoint is identified and incorporated into the coordinated target checkpoint.
Claim 11
The method of claim 10, wherein in response to determining a faulted state for the coordinated target checkpoint due to a particular faulted target storage system, an alternate target storage system that has received a missing local replicated checkpoint is identified and incorporated into the coordinated target checkpoint.
Claim 32
The method of claim 21 further comprising:
applying the coordinated target checkpoint to the replica stored across the plurality of target storage systems.
Claim 12
The method of claim 1 further comprising:
applying the coordinated target checkpoint to the replica dataset stored across the two or more target storage systems.
Claim 33
The method of claim 21 further comprising:
presenting the coordinated target checkpoint as a unified checkpoint for the replica.
Claim 13
The method of claim 1 further comprising:
presenting the coordinated target checkpoint as a unified checkpoint for the replica dataset.
Claim 34
The method of claim 21 further comprising:
coordinating at least one of a clone and a roll back of the replica based on the coordinated target checkpoint.
Claim 14
The method of claim 1 further comprising:
coordinating at least one of a clone and a roll back of the replica dataset based on the coordinated target checkpoint.
Claim 35
The method of claim 21, wherein the coordinated source checkpoint represents a version of the source dataset that does not include any modification that could depend on a result of any other modification not included in the version of the source dataset.
Claim 15
The method of claim 1, wherein the coordinated source checkpoint represents a version of the source dataset that does not include any modification that could depend on a result of any other modification not included in the version of the source dataset.
Although the conflicting claims are not identical, they are not patentably distinct from each other because they are substantially similar in scope and they use the similar limitations to produce the same end result of intelligent managing multi-storage system replication utilizing coordinated dataset checkpoints.
It would have been obvious to a person with ordinary skills in the art at the time of the invention was effectively filed to modify the claims of ‘820 with any combination of the references cited below to arrive at the claims of the instant application for the purpose of enabling large-scale data distribution with high availability and reliability without significantly increasing cost associated with data replication. Further it would have been obvious to modify or to omit the additional elements of claims 1-20 of ‘082 to arrive at claims 21-40 of the instant application because the person would have realized that the remaining element would perform the same functions as before. “Omission of element and its function in combination is obvious expedient if the remaining elements perform same functions as before.” See In re Karlson (CCPA) 136 USPQ 184, decide Jan 16, 1963, Appl. No. 6857, U.S. Court of Customs and Patent Appeals.
Claim Rejections - 35 USC § 101
35 U.S.C. 101 reads as follows:
Whoever invents or discovers any new and useful process, machine, manufacture, or composition of matter, or any new and useful improvement thereof, may obtain a patent therefor, subject to the conditions and requirements of this title.
The claimed invention in claims 21-22, and 24-40 are directed to a judicial exception (i.e., an abstract idea) without significantly more.
Claims 21-22, and 24-40 pass step 1 of the 35 U.S.C. 101 analysis since each claim is either directed to a method, or an apparatus comprising a memory and processing device (i.e., hardware processor and memory per [00228] and [00249] of instant specification).
Claims 21, 36, and 39 recite each, in part, method using steps that are directed to an abstract idea (“Courts have examined claims that required the use of a computer and still found that the underlying, patent-ineligible invention could be performed via pen and paper or in a person’s mind.” Versata Dev. Group v. SAP Am., Inc., 793 F.3d 1306, 1335, 115 USPQ2d 1681, 1702 (Fed. Cir. 2015)). Each claim is directed to the abstract idea of coordinating and replicating data checkpoints across multiple storage systems with a process comprising steps (1) replicating local source checkpoints from source storage systems to target storage systems (e.g., an abstract concept of data duplication wherein moving data from point A to point B is fundamental functional requirement of any backup system), and (2) creating a coordinated target checkpoint for a replica based on those replicated local checkpoints (e.g., an abstract idea of logically ensuring multiple disparate data sets match a master reference). Such process describes the high-level functional goal of organizing and managing data backup/replication which is a longstanding practice in data management considered as “organized human activity.” Each claim focuses on the “what” (i.e., achieving a coordinated state) rather than a specific “how” that improves the underlying hardware. Under the USPTO’s 2019 Revised Patent Subject Matter Eligibility Guidance, this falls under mathematical or mental processes (e.g., logical coordination of timestamps or data states) and certain methods of organizing human activity (e.g., managing storage locations of records and data) since the coordination logic could be performed mentally or manually by a human administrator. That is, other than reciting generic components (e.g., computer processor, computer memory, source/target storage systems), nothing in the claim precludes the limitations from being a process with logics being implemented mentally or manually per step 2A – prong 1 of the “abstract idea” analysis.
In view of step 2A – prong 2 of the “abstract idea” analysis, the claims each does not improve the way a computer operates (e.g., making the processor faster or the memory more efficient). Instead, each claim uses existing computer functions to perform the abstract task of checkpointing. The limitations of each claim are result-oriented as they describe the goal (i.e., a coordinated checkpoint) without reciting any specific algorithm used to achieve that coordination in a way that solves an existing technical problem in distributed computing (e.g., network latency or data race conditions). The limitation are no more than mere instructions to apply the exception using a generic computer component (e.g., processor, memory, and computer-executable instructions).
Each claim, viewed as a whole, under step 2B of the abstract idea analysis, does nothing more than describes the process of managing and synchronizing data across multiple locations using generic computer components. The background of the limitation does not provide any indication that the computer components (e.g., processor, memory, and computer-executable instructions) are not off-the-shelf computer components. The Symantec, TLI, and OOP Techs court decisions cited in MPEP 2106.05(d)(II) indicate that mere receiving, generating, storing, determining, identifying, and transmitting of data over a network are a well-understood, routine, and conventional functions when claimed in a merely generic manner (as it is here). Accordingly, a conclusion that the claims are well-understood, routine, conventional (WURC) activity is supported under Berkheimer Option 2. For these reasons, there is no inventive concept in each claim, thus, claim 21 and 36 are ineligible.
Regarding claim 22, the claim further specifies block storage implementation, file system implementation, database implementation, or object storage implementation for the source/target systems. This merely identifies the field of use or the environment for the abstract idea and does not change the nature of the abstract coordination process. Each specific implementation does not solve an existing technical problem and is WURC implementation to replicate data between different types of storage systems. Thus, the claim is ineligible.
Regarding claim 24, the claim further specifies each local source…system replicates its local…checkpoint to its paired target…system. This describes a basic network topology wherein transmitting data from a source to a target is the fundamental definition of communication. Moving data from one location to another is a WURC activity in the field of backup/replication. Thus, the claim is ineligible.
Regarding claim 25, the claim further specifies at least one replication link employs a replication technology that is different from another replication link. This is a “non-limiting” hardware choice. Merely using different replication technologies does not provide an invention concept because it does not describe how the system overcome any technical difficulty using these different technologies. Thus, the claim is ineligible.
Regarding claim 26, the claim further specifies a source coordinator to facilitate information exchange. Such coordinator is a generic computer functional model being labeled for a functional goal. Simply assigning labels to entities to perform an a WURC abstract task does not add “significantly more” to the claim. Thus, the claim is ineligible.
Regarding claim 27, the claim further specifies a target coordinator to facilitate information exchange. Such coordinator is a generic computer functional model being labeled for a functional goal. Simply assigning labels to entities to perform an a WURC abstract task does not add “significantly more” to the claim. Thus, the claim is ineligible.
Regarding claim 28, the claim merely provide definition for the coordinated source checkpoint as a coordinated snapshot. Thus, the claim is ineligible.
Regarding claim 29, the claim merely provide definitions for a first local replicated checkpoint as a snapshot, and a second local replicated checkpoint as a lightweight checkpoint. These are definitions of data states. Choosing between a full copy (snapshot) and pointer-based copy (lightweight checkpoint) is a routine design choice in data management. These terms describe the nature of the data being stored and do not provide a new way for the computer to function. Thus, the claim is ineligible.
Regarding claims 30, and 37, the claims each further recites determining a state of the coordinated target checkpoint which can be a implemented in a human mind (e.g., mental determination by reading reports of the data states). Thus, the claims are ineligible.
Regarding claim 31, the claim further recites determining a faulted state and identifying an alternate target storage system to incorporate missing checkpoints. This is merely an if-then logical process that can be done in a human mind. Recovering from a failure by looking for an alternative is a basic human problem-solving strategy applied to a computer. It is a functional result of wanting a reliable system ad not a technical solution to a specific hardware failure mechanism. Thus, the claim is ineligible.
Regarding claims 32, 38, and 40, the claims each further recites applying the…checkpoint to the replica… which is a WURC step of writing or saving data in the backup/replication field. Thus, the claims are ineligible.
Regarding claim 33, the claim further recites presenting the…checkpoint… which is a WURC step of displaying output. Thus, the claim is ineligible.
Regarding claim 34, the claim further recites coordinating at least one of a clone and a roll back of the replica which is an extra-solution and WURC activity that does not add an inventive concept to the underlying abstract coordination. Thus, the claim is ineligible.
Regarding claim 35, the claim further recites the…source checkpoint…does not include any modification that could depend on a result of any other modification not included… which is a mathematical/logical algorithm. This limitation is a mental concept of how a set of data should look to be valid. It defines a mathematical relationship between two states of data. Thus, the claim is ineligible.
Claim Rejections - 35 USC § 103
The following is a quotation of 35 U.S.C. 103 which forms the basis for all obviousness rejections set forth in this Office action:
A patent for a claimed invention may not be obtained, notwithstanding that the claimed invention is not identically disclosed as set forth in section 102, if the differences between the claimed invention and the prior art are such that the claimed invention as a whole would have been obvious before the effective filing date of the claimed invention to a person having ordinary skill in the art to which the claimed invention pertains. Patentability shall not be negated by the manner in which the invention was made.
Claims 21, 24-27, and 30-40 are rejected under AIA 35 U.S.C. 103 as being unpatentable over Flynn, JR et al. (Pub. No. US 2007/0244937, published on October 18, 2007; hereinafter Flynn) in view of Annamalai et al. (Pub. No. US 2017/0013058, published on January 12, 2017; hereinafter Annamalai).
Regarding claims 21, 36, and 39, Flynn clearly shows and discloses a method (Abstract); an apparatus comprising a memory; and a processing device operatively coupled to the memory configured to implement the method; and a non-transitory computer readable storage medium storing instructions that, (Figure 4), when executed by a processor, cause a processing device to implement the method comprising:
replicating, from a plurality of source storage systems to a plurality of target storage systems, a plurality of source checkpoints (The data stored in data storage devices A, B, and C at the production site 450 may be transferred to the data storage devices D, E and F at the remotely located recovery site 460 using a peer-to-peer remote copy operation, [0093]-[0095]) that are associated with a coordinated source checkpoint representing a snapshot of a source dataset stored across the plurality of source storage systems (upon initialization of a primary application instance 412 on the primary computing device 410, for example, or at any other suitable time point at which a stateful checkpoint may be generated, the primary fault tolerance engine 414 generates a checkpoint of the state of the primary application instance 412. This checkpoint involves a copy of the application data at the checkpoint time and stateful checkpoint metadata at the checkpoint time. The checkpoint application data copy is generated and stored in the data storage device A while the stateful checkpoint metadata is generated and stored in the data storage device B, [0097]); and
creating, based on the replicated plurality of source checkpoints, a coordinated target checkpoint (The remote fault tolerance engine 424 may initiate a "restart" operation on the shadow application instance 422 in response to the message from the primary fault tolerance engine 414. The restart operation makes use of the copy of the application data and stateful checkpoint metadata to restart the shadow application instance 422 at a state that corresponds to the initial state of the primary application instance 412 specified by the application data and stateful checkpoint metadata, [0099]-[0104]. It is clear that the replica is caused to have the same state and application data as the source systems to constitute a coordinated target checkpoint) for a replica of the source dataset that is stored across the plurality of target storage systems (Figure 4 shows based on the data from data storage devices B and C which are being replicated to storage devices E and F, the remote fault tolerance engine 424 can replay events from both the storage devices E and F to synchronize the shadow application instance 422 with state updates from primary application instance 412).
Annamalai then discloses:
the source checkpoints being local source checkpoints (a first server can be a primary server for a first shard, a secondary server for a second shard and a follower server for a third shard. In some embodiments, the sync replica set can be different for different shards, [0019]. A follower server can commit writes in the exact same order as a primary/secondary does, and maintain the exact same data and replication log as a primary/secondary does, except that the follower server's data get updated with a slight delay compared to that of the primary/secondary, [0034]. It is clear that because each shard is managed by a distinct primary/secondary set and each of those servers maintains its own replication log, the result is a plurality of local logs that collectively form the global dataset);
the replica of the source dataset that is stored across the plurality of target storage systems (Figure 3 shows a distributed computing system such as a social networking application can store data such as user profile data, pictures, messages, comments, etc., associated with users of the social networking application. The data can be partitioned into multiple logical partitions, each of which can be referred to as a shard. For a first shard, the sync replica set can include a first server as a primary server and second and third servers as secondary servers, and for a second shard, the sync replica set can include the first and second servers as the secondary servers and the third server as the primary server. In another example, the sync replica set for the second shard can include a set of servers, e.g., fourth, fifth and sixth servers, that is completely different from that of the first shard, [0019]).
It would have been obvious to an ordinary person skilled in the art at the time of the invention was effectively filed to incorporate the teachings of Annamalai with the teachings of Flynn for the purpose of replicating a distributed dataset from multiple source storages to corresponding destination target storages such that failover can be achieved when a respective storage fails using any available storages maintaining the replicated dataset.
Regarding claim 24, Flynn further discloses each local source storage system replicates its local source checkpoint to its paired target storage system (The active standby computing device 420 is coupled to storage devices D, E and F. Storage devices D, E and F are mirrors of storage devices A, B and C. Thus, storage device D stores application data for the primary application instance 412. Storage device E stores checkpoint metadata for the primary application instance 412 as well as one or more logs of events occurring in application instance 412. Storage device F stores the secondary logs of events occurring in the application instance 412, [0093]-[0095]).
Regarding claim 25, Flynn further discloses at least one replication link employs a replication technology that is different from another replication link (using a combination of asynchronous and synchronous replication methods, [0025]).
Regarding claim 26, Flynn further discloses the coordinated source checkpoint is coordinated across the plurality of source storage systems by a source coordinator (The remote fault tolerance engine 424 may initiate a "restart" operation on the shadow application instance 422 in response to the message from the primary fault tolerance engine 414. The restart operation makes use of the copy of the application data and stateful checkpoint metadata to restart the shadow application instance 422 at a state that corresponds to the initial state of the primary application instance 412 specified by the application data and stateful checkpoint metadata, [0099]-[0104]. Figure 4 shows based on the data from data storage devices B and C which are being replicated to storage devices E and F, the remote fault tolerance engine 424 can replay events from both the storage devices E and F to synchronize the shadow application instance 422 with state updates from primary application instance 412).
Regarding claim 27, Flynn further discloses a target coordinator receives, from the source coordinator, information relating to the coordinated source checkpoint (The remote fault tolerance engine 424 may initiate a "restart" operation on the shadow application instance 422 in response to the message from the primary fault tolerance engine 414. The restart operation makes use of the copy of the application data and stateful checkpoint metadata to restart the shadow application instance 422 at a state that corresponds to the initial state of the primary application instance 412 specified by the application data and stateful checkpoint metadata, [0099]-[0104]. Figure 4 shows based on the data from data storage devices B and C which are being replicated to storage devices E and F, the remote fault tolerance engine 424 can replay events from both the storage devices E and F to synchronize the shadow application instance 422 with state updates from primary application instance 412).
Regarding claims 30, and 37, Flynn further discloses wherein creating, based on two or more local replicated checkpoints, a coordinated target checkpoint for the replica includes: determining a state of the coordinated target checkpoint (The remote fault tolerance engine 424 may initiate a "restart" operation on the shadow application instance 422 in response to the message from the primary fault tolerance engine 414. The restart operation makes use of the copy of the application data and stateful checkpoint metadata to restart the shadow application instance 422 at a state that corresponds to the initial state of the primary application instance 412 specified by the application data and stateful checkpoint metadata, [0099]-[0104]. Figure 4 shows based on the data from data storage devices B and C which are being replicated to storage devices E and F, the remote fault tolerance engine 424 can replay events from both the storage devices E and F to synchronize the shadow application instance 422 with state updates from primary application instance 412).
Regarding claim 31, Annamalai then discloses in response to determining a faulted state for the coordinated target checkpoint due to a particular faulted target storage system, an alternate target storage system that has received a missing local replicated checkpoint is identified and incorporated into the coordinated target checkpoint (the deployment topology may not have to define that a specific server computer has to be in a specific hierarchical level. In some embodiments, this can give flexibility in deploying any server in any hierarchical level and/or for any shard. Further, if more number of servers become available over time, additional servers can be deployed in any of the hierarchical levels automatically, or if a specified server fails in a specified hierarchical level, another server can be deployed in the specified hierarchical level, [0023]. The sync replica set 105 is an example of a “3-way” replication, in which the distributed computing system 150 can guarantee three replicas of the data. Typically, in a “3-way” replication, a sync replica set includes one primary server and two secondary servers. The number of servers in a sync replica set is configurable and can depend on various factors, e.g., fault tolerance, load balancing, write latency and other required performance characteristics, [0028]).
Regarding claims 32, 38, and 40, Flynn further discloses applying the coordinated target checkpoint to the replica dataset stored across the plurality of target storage systems (The remote fault tolerance engine 424 may initiate a "restart" operation on the shadow application instance 422 in response to the message from the primary fault tolerance engine 414. The restart operation makes use of the copy of the application data and stateful checkpoint metadata to restart the shadow application instance 422 at a state that corresponds to the initial state of the primary application instance 412 specified by the application data and stateful checkpoint metadata, [0099]-[0104]. Figure 4 shows based on the data from data storage devices B and C which are being replicated to storage devices E and F, the remote fault tolerance engine 424 can replay events from both the storage devices E and F to synchronize the shadow application instance 422 with state updates from primary application instance 412).
Regarding claim 33, Flynn further discloses presenting the coordinated target checkpoint as a unified checkpoint for the replica (Figure 4 shows data stored in storage devices D-F are combined to be used by remote fault tolerance engine 424 for state synchronization with primary fault tolerance engine 414, [0093]-[0095]).
Regarding claim 34, Flynn further discloses coordinating at least one of a clone and a roll back of the replica based on the coordinated target checkpoint (The remote fault tolerance engine 424 may initiate a "restart" operation on the shadow application instance 422 in response to the message from the primary fault tolerance engine 414. The restart operation makes use of the copy of the application data and stateful checkpoint metadata to restart the shadow application instance 422 at a state that corresponds to the initial state of the primary application instance 412 specified by the application data and stateful checkpoint metadata, [0099]-[0104]. Figure 4 shows based on the data from data storage devices B and C which are being replicated to storage devices E and F, the remote fault tolerance engine 424 can replay events from both the storage devices E and F to synchronize the shadow application instance 422 with state updates from primary application instance 412).
Regarding claim 35, Flynn further discloses the coordinated source checkpoint represents a version of the source dataset that does not include any modification that could depend on a result of any other modification not included in the version of the source dataset (at any other suitable time point at which a stateful checkpoint may be generated, the primary fault tolerance engine 414 generates a checkpoint of the state of the primary application instance 412. This checkpoint involves a copy of the application data at the checkpoint time and stateful checkpoint metadata at the checkpoint time, [0097]).
Claim 28 is rejected under AIA 35 U.S.C. 103 as being unpatentable over Flynn in view of Annamalai and further in view of Yin et al. (Pub. No. US 2020/0250151, filed on January 31, 2019; hereinafter Yin).
Regarding claim 28, Yin then discloses the plurality of source storage systems store respective distinct local portion of the source dataset (the source storage platform 302 may be utilized to store a database 314 comprised of records 316. The records 316 are distributed over one or more replication sets. Each replication set (e.g., shard) includes a distinct data set. For example, the database 314 may include three shards including shard “A” 318 comprised of records of employees of XYZ Corp. with last names beginning with letters A-J, shard “B” 320, comprised of records of employees of XYZ Corp. with last names beginning with letters K-T, shard “C” 322, comprised of records of employees of XYZ Corp. with last names beginning with letters U-Z, [0002], [0087]).
It would have been obvious to an ordinary person skilled in the art at the time of the invention was effectively filed to incorporate the teachings of Yin with the teachings of Flynn, as modified by Annamalai, for the purpose of ensuring point-int-time state consistency across distributed dataset shards based on operation logs associated with the storages associated with the distributed dataset.
Allowable Subject Matter
Claims 22, and 29 are objected for being dependent on a base rejected claim but would be allowable (given that the 35 U.S.C. 101 rejections are resolved) if rewritten in independent form to incorporate the limitation of the base claim and all intervening claim(s).
Claim 23 is objected for being dependent on a base rejected claim but would be allowable if rewritten in independent form to incorporate the limitation of the base claim and all intervening claim(s).
Relevant Prior Art
The following prior art is/are deemed relevant to the claims:
Olston et al. (Pub. No. US 2009/0307329) teaches if a task is scheduled on a machine that does not yet have copies of the portions of the data set on which the task needs to operate, then that machine obtains copies of those portions from other machines that already have those portions. According to one embodiment of the invention, whenever a "source" machine ships a portion of a data set to another "destination" machine in the distributed system, the destination machine makes a persistent, local copy of that portion on the destination machine's persistent storage mechanism. Thus, portions of the data set are automatically replicated whenever those portions are shipped between machines of the distributed system.
Conclusion
THIS ACTION IS MADE FINAL. Applicant is reminded of the extension of time policy as set forth in 37 CFR 1.136(a).
A shortened statutory period for reply to this final action is set to expire THREE MONTHS from the mailing date of this action. In the event a first reply is filed within TWO MONTHS of the mailing date of this final action and the advisory action is not mailed until after the end of the THREE-MONTH shortened statutory period, then the shortened statutory period will expire on the date the advisory action is mailed, and any nonprovisional extension fee (37 CFR 1.17(a)) pursuant to 37 CFR 1.136(a) will be calculated from the mailing date of the advisory action. In no event, however, will the statutory period for reply expire later than SIX MONTHS from the mailing date of this final action.
Contact Information
Any inquiry concerning this communication or earlier communications from the Examiner should be directed to Son Hoang whose telephone number is (571) 270-1752. The Examiner can normally be reached on Monday – Friday (7:00 AM – 4:00 PM).
If attempts to reach the Examiner by telephone are unsuccessful, the Examiner’s supervisor, Sherief Badawi can be reached on (571) 272-9782. The fax phone number for the organization where this application or proceeding is assigned is 571-273-8300.
Information regarding the status of an application may be obtained from the Patent Application Information Retrieval (PAIR) system. Status information for published applications may be obtained from either Private PAIR or Public PAIR. Status information for unpublished applications is available through Private PAIR only. For more information about the PAIR system, see http://pair-direct.uspto.gov. Should you have questions on access to the Private PAIR system, contact the Electronic Business Center (EBC) at 866-217-9197 (toll-free). If you would like assistance from a USPTO Customer Service Representative or access to the automated information system, call 800-786-9199 (IN USA OR CANADA) or 571-272-1000.
/SON T HOANG/
Primary Examiner, Art Unit 2169 August 6, 2026